[SUPPORT] Failed Flink Checkpoints
- Dominant language
- Java
- Stars
- 6.2k
- Forks
- 2.5k
- Avg merge
- 2d 8h
- Merged PRs (30d)
- 111
Description
**_Tips before filing an issue_**
- Have you gone through our [FAQs](https://hudi.apache.org/learn/faq/)?
- Join the mailing list to engage in conversations and get faster support at dev-subscribe@hudi.apache.org.
- If you have triaged this as a bug, then file an [issue](https://issues.apache.org/jira/projects/HUDI/issues) directly.
**Describe the problem you faced**
A clear and concise description of the problem.
**To Reproduce**
Steps to reproduce the behavior:
1.
2.
3.
4.
**Expected behavior**
We expect successful checkpoints using Hudi with Flink
**Environment Description**
Using AWS Managed Flink Applications.
Hudi table config:
```
'table.type' = 'COPY_ON_WRITE',
'hoodie.datasource.write.recordkey.field' = 'event_id',
'hoodie.datasource.write.partitionpath.field' = 'customer_uuid,year,month',
'hoodie.datasource.write.precombine.field' = 'time',
'hoodie.schema.on.read.enable' = 'true',
'hoodie.compaction.inline' = 'false',
'hoodie.parquet.compression.codec' = 'snappy',
'compaction.async.enabled' = 'false',
'hoodie.datasource.write.operation' = 'upsert',
'compaction.tasks' = '0',
'write.tasks' = '8',
'hoodie.populate.meta.fields' = 'false',
'hoodie.table.keygenerator.class' = 'org.apache.hudi.keygen.CustomKeyGenerator',
'state.backend.incremental' = 'true',
'write.merge.max_memory' = '1024',
'write.task.max.size' = '2014D',
'state.backend' = 'rocksdb',
'hoodie.metadata.enable' = 'false'
```
Flink encounters issues when generating checkpoints for hudi stream_write operator:
Exception raised:
```
Caused by: org.apache.flink.runtime.checkpoint.CheckpointException: Could not complete snapshot 4 for operator stream_write: default_database.hudi_access_logs (2/8)#2. Failure reason: Checkpoint was declined.
at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:300)
at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:204)
at org.apache.flink.streaming.api.operators.AbstractStreamOperator.snapshotState(AbstractStreamOperator.java:395)
at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.checkpointStreamOperator(RegularOperatorChain.java:228)
at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.buildOperatorSnapshotFutures(RegularOperatorChain.java:213)
at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.snapshotState(RegularOperatorChain.java:192)
at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.takeSnapshotSync(SubtaskCheckpointCoordinatorImpl.java:751)
at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointState(SubtaskCheckpointCoordinatorImpl.java:361)
at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$18(StreamTask.java:1437)
at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50)
at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:1425)
at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:1382)
... 22 more
Caused by: java.lang.IllegalArgumentException
at org.apache.hudi.common.util.ValidationUtils.checkArgument(ValidationUtils.java:33)
at org.apache.hudi.io.HoodieMergeHandle.validateAndSetAndKeyGenProps(HoodieMergeHandle.java:152)
at org.apache.hudi.io.HoodieMergeHandle.(HoodieMergeHandle.java:135)
at org.apache.hudi.io.HoodieMergeHandle.(HoodieMergeHandle.java:125)
at org.apache.hudi.io.FlinkMergeHandle.(FlinkMergeHandle.java:69)
at org.apache.hudi.io.FlinkWriteHandleFactory$CommitWriteHandleFactory.createMergeHandle(FlinkWriteHandleFactory.java:181)
at org.apache.hudi.io.FlinkWriteHandleFactory$BaseCommitWriteHandleFactory.create(FlinkWriteHandleFactory.java:123)
at org.apache.hudi.client.HoodieFlinkWriteClient.getOrCreateWriteHandle(HoodieFlinkWriteClient.java:459)
at org.apache.hudi.client.HoodieFlinkWriteClient.access$000(HoodieFlinkWriteClient.java:77)
at org.apache.hudi.client.HoodieFlinkWriteClient$AutoCloseableWriteHandle.(HoodieFlinkWriteClient.java:515)
at org.apache.hudi.client.HoodieFlinkWriteClient$AutoCloseableWriteHandle.(HoodieFlinkWriteClient.java:507)
at org.apache.hudi.client.HoodieFlinkWriteClient.upsert(HoodieFlinkWriteClient.java:148)
at org.apache.hudi.sink.StreamWriteFunction.lambda$initWriteFunction$1(StreamWriteFunction.java:192)
at org.apache.hudi.sink.StreamWriteFunction.writeBucket(StreamWriteFunction.java:495)
at org.apache.hudi.sink.StreamWriteFunction.lambda$flushRemaining$7(StreamWriteFunction.java:467)
at java.base/java.util.LinkedHashMap$LinkedValues.forEach(LinkedHashMap.java:608)
at org.apache.hudi.sink.StreamWriteFunction.flushRemaining(StreamWriteFunction.java:463)
at org.apache.hudi.sink.StreamWriteFunction.snapshotState(StreamWriteFunction.java:137)
at org.apache.hudi.sink.common.AbstractStreamWriteFunction.snapshotState(AbstractStreamWriteFunction.java:167)
at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:118)
at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:99)
at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.snapshotState(AbstractUdfStreamOperator.java:89)
at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:253)
... 33 more
```
* Hudi version : 0.15.0
* Flink version : 1.20
* Storage (HDFS/S3/GCS..) : S3
**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 failure with Hudi 0.15.0, Flink 1.20, S3, and the listed table configuration. Read HoodieMergeHandle.validateAndSetAndKeyGenProps, FlinkWriteHandleFactory, and StreamWriteFunction.snapshotState; done means identifying the invalid condition and demonstrating successful checkpoints for the stream_write operator.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- aws, java
- Domain
- cloud, stream-processing
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100