apache / apache/hudi

[SUPPORT] Flink Incremental read task use 'payload.class' configure does not work

Open
#10,351 1 comment 0 reactions 0 assignees View on GitHub
engine:flink priority:medium
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

**Describe the problem you faced**

In Flink task,I setup payload.class as customPayload, but it does not work.

How should I configure custom payload in flink read task instead of change flink write task, and does not change hoodie.compaction.payload.class configure?
just only change read task payload logical.

**To Reproduce**

create Flink sql table:
````
CREATE TEMPORARY TABLE IF NOT EXISTS t1
(
`id` BIGINT,
`name` STRING,
PRIMARY KEY(id) NOT ENFORCED
)
WITH(
'hoodie.datasource.write.payload.class'='xxxx.CustomPayload',
'payload.class'='xxxx.CustomPayload',
'metadata.enabled'='false',
'compaction.async.enabled'='false',
'hoodie.datasource.query.type'='incremental',
'connector'='hudi',
'index.type'='BUCKET',
'hoodie.bucket.index.num.buckets'='8',
'read.data.skipping.enabled'='true',
'read.streaming.enabled'='false',
'path'='hdfs://xxxx/path',
'read.tasks'='4',
'read.start-commit'='20231212004636502',
'read.end-commit'='20231218001637223',
'table.type'='MERGE_ON_READ',
'compaction.schedule.enabled'='false',
'changelog.enabled'='true',
'hoodie.bucket.index.hash.field'='id'
)
````

hoodie.properties:
````
hoodie.table.precombine.field=_time
hoodie.datasource.write.drop.partition.columns=false
hoodie.table.type=MERGE_ON_READ
hoodie.archivelog.folder=archived
hoodie.table.cdc.enabled=false
hoodie.compaction.payload.class=org.apache.hudi.common.model.EventTimeAvroPayload
hoodie.timeline.layout.version=1
hoodie.table.version=5
hoodie.table.recordkey.fields=id
hoodie.datasource.write.partitionpath.urlencode=false
hoodie.table.name=t1
hoodie.table.keygenerator.class=org.apache.hudi.keygen.NonpartitionedAvroKeyGenerator
hoodie.compaction.record.merger.strategy=eeb8d96f-b1e4-49fd-bbf8-28ac514178e5
....
````

**Expected behavior**

maybe it should use custom payload merge record and output, instead of use EventTimeAvroPayload.

**Environment Description**

* Hudi version : 0.13.1

* Flink version : 1.15.4

* Hadoop version : 3.0.0

* Running on Docker? (yes/no) : no

**Additional context**

Add any other context about the problem here.

**Stacktrace**

```Add the stacktrace of the error.```

Contributor guide

No contributing guide indexed for this repository

Research direction

Start by reproducing the Flink incremental read task with the provided SQL table options and hoodie.properties on Hudi 0.13.1, Flink 1.15.4, and Hadoop 3.0. Compare the read-side payload.class behavior with hoodie.compaction.payload.class. Done means the read task uses the configured custom payload for merging and output without changing the write or compaction payload configuration.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.