apache / apache/dolphinscheduler

[Bug] [Worker Error] The heartbeat request received by the Flink REST API is in an incorrect format and cannot be deserialized

Open
#17,321 6 comments 0 reactions 1 assignee Claimed by @zhoujianpa View on GitHub
bug help wanted
Dominant language
Java
Stars
14.5k
Forks
5.1k
Avg merge
1d 21h
Merged PRs (30d)
29

Description

### Search before asking

- [x] I had searched in the [issues](https://github.com/apache/dolphinscheduler/issues?q=is%3Aissue) and found no similar issues.

### What happened

DolphinScheduler calls flink and is actively canceled

### What you expected to happen

Flink is running normally

### How to reproduce

flink 1.18.0 standlone
DolphinScheduler 3.2.0 Pseudo-Cluste
jdk 1.8
flink java mysqlCdc synchronization task local execution ,DolphinScheduler calls flink 4 minutes later.
I think the interface content of DolphinScheduler and flink's REST api interaction is inconsistent. Mount the jar package to flink and execute it normally.

### Anything else

flink error log:
org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Marking checkpoint 26 as completed for source Source: MySQL Source.
2025-07-03 10:39:27,792 DEBUG org.apache.flink.runtime.jobmaster.JobMaster [] - Trigger heartbeat request.
2025-07-03 10:39:27,795 DEBUG org.apache.flink.runtime.jobmaster.JobMaster [] - Received heartbeat from localhost:44787-15170a.
2025-07-03 10:39:29,213 INFO org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Triggering checkpoint 27 (type=CheckpointType{name='Checkpoint', sharingFilesStrategy=FORWARD_BACKWARD}) @ 1751510369213 for job e6b7b65a156d4872c755189aecf3147e.
2025-07-03 10:39:29,214 DEBUG org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Taking a state snapshot on operator Source: MySQL Source for checkpoint 27
2025-07-03 10:39:29,232 DEBUG org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Received acknowledge message for checkpoint 27 from task 4ec101375946135b749b21299edb2531_d618a97df21bbd4bb61c79cdeca965b4_0_0 of job e6b7b65a156d4872c755189aecf3147e at localhost:44787-15170a @ localhost (dataPort=43623).
2025-07-03 10:39:29,232 DEBUG org.apache.flink.runtime.state.SharedStateRegistryImpl [] - Discard state created before checkpoint 27 and not used afterwards
2025-07-03 10:39:29,232 DEBUG org.apache.flink.runtime.checkpoint.CheckpointsCleaner [] - Try to discard checkpoint 26.
2025-07-03 10:39:29,233 DEBUG org.apache.flink.runtime.checkpoint.CheckpointStatsTracker [] - CheckpointStatistics (for jobID=e6b7b65a156d4872c755189aecf3147e, checkpointId=27) dump = {"className":"completed","id":27,"status":"COMPLETED","is_savepoint":false,"savepointFormat":null,"trigger_timestamp":1751510369213,"latest_ack_timestamp":1751510369232,"checkpointed_size":2877,"state_size":2877,"end_to_end_duration":19,"alignment_buffered":0,"processed_data":0,"persisted_data":0,"num_subtasks":1,"num_acknowledged_subtasks":1,"checkpoint_type":"CHECKPOINT","tasks":{"d618a97df21bbd4bb61c79cdeca965b4":{"id":27,"status":"COMPLETED","latest_ack_timestamp":1751510369232,"checkpointed_size":2877,"state_size":2877,"end_to_end_duration":19,"alignment_buffered":0,"processed_data":0,"persisted_data":0,"num_subtasks":1,"num_acknowledged_subtasks":1}},"external_path":"file:/mnt/vortex/baseserver/flink/data/checkpoints/e6b7b65a156d4872c755189aecf3147e/chk-27","discarded":false}
2025-07-03 10:39:29,233 INFO org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Completed checkpoint 27 for job e6b7b65a156d4872c755189aecf3147e (5408 bytes, checkpointDuration=19 ms, finalizationTime=1 ms).
2025-07-03 10:39:29,233 DEBUG org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Checkpoint state: OperatorState(operatorID: d618a97df21bbd4bb61c79cdeca965b4, parallelism: 1, maxParallelism: 128, coordinatorState: 2531 bytes, sub task states: 1, total size (bytes): 5408), OperatorState(operatorID: 3b579097cd35da2657ad2fdca381c2f8, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 3800542a0b08d84a487282e017c42b42, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 602a61369b9434e3f66ad4f48ab891cb, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: a34dbf8336b99406273a3eae87fcecff, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 8bc47f9a3498f52aa521a5f46d88bdef, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0)
2025-07-03 10:39:29,233 INFO org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Marking checkpoint 27 as completed for source Source: MySQL Source.
2025-07-03 10:39:32,213 INFO org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Triggering checkpoint 28 (type=CheckpointType{name='Checkpoint', sharingFilesStrategy=FORWARD_BACKWARD}) @ 1751510372213 for job e6b7b65a156d4872c755189aecf3147e.
2025-07-03 10:39:32,214 DEBUG org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Taking a state snapshot on operator Source: MySQL Source for checkpoint 28
2025-07-03 10:39:32,220 DEBUG org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Received acknowledge message for checkpoint 28 from task 4ec101375946135b749b21299edb2531_d618a97df21bbd4bb61c79cdeca965b4_0_0 of job e6b7b65a156d4872c755189aecf3147e at localhost:44787-15170a @ localhost (dataPort=43623).
2025-07-03 10:39:32,221 DEBUG org.apache.flink.runtime.state.SharedStateRegistryImpl [] - Discard state created before checkpoint 28 and not used afterwards
2025-07-03 10:39:32,221 DEBUG org.apache.flink.runtime.checkpoint.CheckpointsCleaner [] - Try to discard checkpoint 27.
2025-07-03 10:39:32,221 DEBUG org.apache.flink.runtime.checkpoint.CheckpointStatsTracker [] - CheckpointStatistics (for jobID=e6b7b65a156d4872c755189aecf3147e, checkpointId=28) dump = {"className":"completed","id":28,"status":"COMPLETED","is_savepoint":false,"savepointFormat":null,"trigger_timestamp":1751510372213,"latest_ack_timestamp":1751510372220,"checkpointed_size":2877,"state_size":2877,"end_to_end_duration":7,"alignment_buffered":0,"processed_data":0,"persisted_data":0,"num_subtasks":1,"num_acknowledged_subtasks":1,"checkpoint_type":"CHECKPOINT","tasks":{"d618a97df21bbd4bb61c79cdeca965b4":{"id":28,"status":"COMPLETED","latest_ack_timestamp":1751510372220,"checkpointed_size":2877,"state_size":2877,"end_to_end_duration":7,"alignment_buffered":0,"processed_data":0,"persisted_data":0,"num_subtasks":1,"num_acknowledged_subtasks":1}},"external_path":"file:/mnt/vortex/baseserver/flink/data/checkpoints/e6b7b65a156d4872c755189aecf3147e/chk-28","discarded":false}
2025-07-03 10:39:32,221 INFO org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Completed checkpoint 28 for job e6b7b65a156d4872c755189aecf3147e (5408 bytes, checkpointDuration=8 ms, finalizationTime=0 ms).
2025-07-03 10:39:32,221 DEBUG org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Checkpoint state: OperatorState(operatorID: d618a97df21bbd4bb61c79cdeca965b4, parallelism: 1, maxParallelism: 128, coordinatorState: 2531 bytes, sub task states: 1, total size (bytes): 5408), OperatorState(operatorID: 3b579097cd35da2657ad2fdca381c2f8, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 3800542a0b08d84a487282e017c42b42, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 602a61369b9434e3f66ad4f48ab891cb, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: a34dbf8336b99406273a3eae87fcecff, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 8bc47f9a3498f52aa521a5f46d88bdef, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0)
2025-07-03 10:39:32,222 INFO org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Marking checkpoint 28 as completed for source Source: MySQL Source.
2025-07-03 10:39:34,980 DEBUG org.apache.flink.runtime.resourcemanager.StandaloneResourceManager [] - Trigger heartbeat request.
2025-07-03 10:39:34,980 DEBUG org.apache.flink.runtime.resourcemanager.StandaloneResourceManager [] - Trigger heartbeat request.
2025-07-03 10:39:34,980 DEBUG org.apache.flink.runtime.jobmaster.JobMaster [] - Received heartbeat request from a14492ad199fe51ed88127acd6f5d40e.
2025-07-03 10:39:34,980 DEBUG org.apache.flink.runtime.resourcemanager.StandaloneResourceManager [] - Received heartbeat from 67c23a2c159ae2df96669e309d401504.
2025-07-03 10:39:34,983 DEBUG org.apache.flink.runtime.resourcemanager.StandaloneResourceManager [] - Received heartbeat from localhost:44787-15170a.
2025-07-03 10:39:34,983 DEBUG org.apache.flink.runtime.resourcemanager.slotmanager.FineGrainedSlotManager [] - Received slot report from instance 64b49f30e03aea0e11ea1ea9250b6740: SlotReport{
SlotStatus{slotID=localhost:44787-15170a_0, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_1, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_2, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_3, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_4, allocationID=845383d5444cb0613257c2c7a0221485, jobID=e6b7b65a156d4872c755189aecf3147e, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}}.
2025-07-03 10:39:34,984 DEBUG org.apache.flink.runtime.resourcemanager.slotmanager.DefaultSlotStatusSyncer [] - Received slot report from instance 64b49f30e03aea0e11ea1ea9250b6740: SlotReport{
SlotStatus{slotID=localhost:44787-15170a_0, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_1, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_2, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_3, allocationID=null, jobID=null, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}
SlotStatus{slotID=localhost:44787-15170a_4, allocationID=845383d5444cb0613257c2c7a0221485, jobID=e6b7b65a156d4872c755189aecf3147e, resourceProfile=ResourceProfile{cpuCores=1, taskHeapMemory=364.800mb (382520517 bytes), taskOffHeapMemory=0 bytes, managedMemory=343.040mb (359703515 bytes), networkMemory=85.760mb (89925878 bytes)}}}.
2025-07-03 10:39:34,984 DEBUG org.apache.flink.runtime.io.network.partition.ResourceManagerPartitionTrackerImpl [] - Processing cluster partition report from task executor localhost:44787-15170a: PartitionReport{entries=[]}.
2025-07-03 10:39:35,213 INFO org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Triggering checkpoint 29 (type=CheckpointType{name='Checkpoint', sharingFilesStrategy=FORWARD_BACKWARD}) @ 1751510375213 for job e6b7b65a156d4872c755189aecf3147e.
2025-07-03 10:39:35,214 DEBUG org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Taking a state snapshot on operator Source: MySQL Source for checkpoint 29
2025-07-03 10:39:35,220 DEBUG org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Received acknowledge message for checkpoint 29 from task 4ec101375946135b749b21299edb2531_d618a97df21bbd4bb61c79cdeca965b4_0_0 of job e6b7b65a156d4872c755189aecf3147e at localhost:44787-15170a @ localhost (dataPort=43623).
2025-07-03 10:39:35,221 DEBUG org.apache.flink.runtime.state.SharedStateRegistryImpl [] - Discard state created before checkpoint 29 and not used afterwards
2025-07-03 10:39:35,221 DEBUG org.apache.flink.runtime.checkpoint.CheckpointsCleaner [] - Try to discard checkpoint 28.
2025-07-03 10:39:35,221 DEBUG org.apache.flink.runtime.checkpoint.CheckpointStatsTracker [] - CheckpointStatistics (for jobID=e6b7b65a156d4872c755189aecf3147e, checkpointId=29) dump = {"className":"completed","id":29,"status":"COMPLETED","is_savepoint":false,"savepointFormat":null,"trigger_timestamp":1751510375213,"latest_ack_timestamp":1751510375220,"checkpointed_size":2877,"state_size":2877,"end_to_end_duration":7,"alignment_buffered":0,"processed_data":0,"persisted_data":0,"num_subtasks":1,"num_acknowledged_subtasks":1,"checkpoint_type":"CHECKPOINT","tasks":{"d618a97df21bbd4bb61c79cdeca965b4":{"id":29,"status":"COMPLETED","latest_ack_timestamp":1751510375220,"checkpointed_size":2877,"state_size":2877,"end_to_end_duration":7,"alignment_buffered":0,"processed_data":0,"persisted_data":0,"num_subtasks":1,"num_acknowledged_subtasks":1}},"external_path":"file:/mnt/vortex/baseserver/flink/data/checkpoints/e6b7b65a156d4872c755189aecf3147e/chk-29","discarded":false}
2025-07-03 10:39:35,221 INFO org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Completed checkpoint 29 for job e6b7b65a156d4872c755189aecf3147e (5408 bytes, checkpointDuration=8 ms, finalizationTime=0 ms).
2025-07-03 10:39:35,221 DEBUG org.apache.flink.runtime.checkpoint.CheckpointCoordinator [] - Checkpoint state: OperatorState(operatorID: d618a97df21bbd4bb61c79cdeca965b4, parallelism: 1, maxParallelism: 128, coordinatorState: 2531 bytes, sub task states: 1, total size (bytes): 5408), OperatorState(operatorID: 3b579097cd35da2657ad2fdca381c2f8, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 3800542a0b08d84a487282e017c42b42, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 602a61369b9434e3f66ad4f48ab891cb, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: a34dbf8336b99406273a3eae87fcecff, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0), OperatorState(operatorID: 8bc47f9a3498f52aa521a5f46d88bdef, parallelism: 1, maxParallelism: 128, coordinatorState: (none), sub task states: 1, total size (bytes): 0)
2025-07-03 10:39:35,221 INFO org.apache.flink.runtime.source.coordinator.SourceCoordinator [] - Marking checkpoint 29 as completed for source Source: MySQL Source.
2025-07-03 10:39:37,792 DEBUG org.apache.flink.runtime.jobmaster.JobMaster [] - Trigger heartbeat request.
2025-07-03 10:39:37,796 DEBUG org.apache.flink.runtime.jobmaster.JobMaster [] - Received heartbeat from localhost:44787-15170a.
2025-07-03 10:39:37,802 ERROR org.apache.flink.runtime.rest.handler.job.JobClientHeartbeatHandler [] - Exception occurred in REST handler.
org.apache.flink.runtime.rest.handler.RestHandlerException: Request did not match expected format JobClientHeartbeatRequestBody.
at org.apache.flink.runtime.rest.handler.AbstractHandler.respondAsLeader(AbstractHandler.java:165) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.LeaderRetrievalHandler.lambda$channelRead0$0(LeaderRetrievalHandler.java:83) ~[flink-dist-1.18.0.jar:1.18.0]
at java.util.Optional.ifPresent(Optional.java:178) ~[?:?]
at org.apache.flink.util.OptionalConsumer.ifPresent(OptionalConsumer.java:45) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.LeaderRetrievalHandler.channelRead0(LeaderRetrievalHandler.java:80) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.LeaderRetrievalHandler.channelRead0(LeaderRetrievalHandler.java:49) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:99) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.router.RouterHandler.routed(RouterHandler.java:115) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.router.RouterHandler.channelRead0(RouterHandler.java:94) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.router.RouterHandler.channelRead0(RouterHandler.java:55) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:99) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.handler.codec.MessageToMessageDecoder.channelRead(MessageToMessageDecoder.java:103) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.FileUploadHandler.channelRead0(FileUploadHandler.java:208) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.FileUploadHandler.channelRead0(FileUploadHandler.java:69) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:99) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.CombinedChannelDuplexHandler$DelegatingChannelHandlerContext.fireChannelRead(CombinedChannelDuplexHandler.java:436) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:346) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:318) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.CombinedChannelDuplexHandler.channelRead(CombinedChannelDuplexHandler.java:251) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:442) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:440) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:166) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:788) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:724) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:650) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:562) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:997) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.netty4.io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) ~[flink-dist-1.18.0.jar:1.18.0]
at java.lang.Thread.run(Thread.java:1583) [?:?]
Caused by: org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.MismatchedInputException: Cannot construct instance of org.apache.flink.runtime.rest.messages.JobClientHeartbeatRequestBody (although at least one Creator exists): cannot deserialize from Object value (no delegate- or property-based Creator)
at [Source: (org.apache.flink.shaded.netty4.io.netty.buffer.ByteBufInputStream); line: 1, column: 2]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.MismatchedInputException.from(MismatchedInputException.java:63) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.DeserializationContext.reportInputMismatch(DeserializationContext.java:1733) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.DeserializationContext.handleMissingInstantiator(DeserializationContext.java:1358) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializerBase.deserializeFromObjectUsingNonDefault(BeanDeserializerBase.java:1420) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializer.deserializeFromObject(BeanDeserializer.java:352) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializer.deserialize(BeanDeserializer.java:185) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.DefaultDeserializationContext.readRootValue(DefaultDeserializationContext.java:323) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper._readMapAndClose(ObjectMapper.java:4730) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:3714) ~[flink-dist-1.18.0.jar:1.18.0]
at org.apache.flink.runtime.rest.handler.AbstractHandler.respondAsLeader(AbstractHandler.java:162) ~[flink-dist-1.18.0.jar:1.18.0]
... 45 more

dolphinscheduler worker erro log:
[INFO] 2025-07-03 10:39:23.071 +0800 org.apache.dolphinscheduler.server.worker.message.MessageRetryRunner:[129] - [WorkflowInstance-0][TaskInstance-40] - Success send message to master, message: TaskExecuteResultMessage(super=BaseMessage(messageSenderAddress=10.10.10.196:1234, messageReceiverAddress=10.10.10.196:5678, messageSendTime=1751510363070), taskInstanceId=40, processInstanceId=34, status=9, startTime=1751446362591, host=10.10.10.196:1234, logPath=/mnt/apache-dolphinscheduler-3.2.0-bin/worker-server/logs/20250702/18178025745248/2/34/40.log, executePath=/tmp/dolphinscheduler/exec/process/default/18164315973344/18178025745248_2/34/40, endTime=1751446438703, processId=762707, appIds=, varPool=[])
[INFO] 2025-07-03 10:42:09.669 +0800 org.apache.dolphinscheduler.plugin.task.api.AbstractTask:[181] - [WorkflowInstance-0][TaskInstance-0] - ->
java.util.concurrent.ExecutionException: org.apache.flink.client.program.ProgramInvocationException: Job failed (JobID: e6b7b65a156d4872c755189aecf3147e)
at java.base/java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:396)
at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2073)
at org.apache.flink.client.program.StreamContextEnvironment.getJobExecutionResult(StreamContextEnvironment.java:171)
at org.apache.flink.client.program.StreamContextEnvironment.execute(StreamContextEnvironment.java:122)
at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:2099)
at org.example.FlinkCdc.main(FlinkCdc.java:199)
at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(DirectMethodHandleAccessor.java:103)
at java.base/java.lang.reflect.Method.invoke(Method.java:580)
at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:355)
at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222)
at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:105)
at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:851)
at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:245)
at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1095)
at org.apache.flink.client.cli.CliFrontend.lambda$mainInternal$9(CliFrontend.java:1189)
at java.base/java.security.AccessController.doPrivileged(AccessController.java:714)
at java.base/javax.security.auth.Subject.doAs(Subject.java:525)
at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1836)
at org.apache.flink.runtime.security.contexts.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
at org.apache.flink.client.cli.CliFrontend.mainInternal(CliFrontend.java:1189)
at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1157)
Caused by: org.apache.flink.client.program.ProgramInvocationException: Job failed (JobID: e6b7b65a156d4872c755189aecf3147e)
at org.apache.flink.client.deployment.ClusterClientJobClientAdapter.lambda$null$6(ClusterClientJobClientAdapter.java:130)
at java.base/java.util.concurrent.CompletableFuture$UniApply.tryFire(CompletableFuture.java:646)
at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:510)
at java.base/java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:2179)
at org.apache.flink.util.concurrent.FutureUtils.lambda$retryOperationWithDelay$6(FutureUtils.java:302)
at java.base/java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:863)
at java.base/java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:841)
at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:510)
at java.base/java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:2179)
at org.apache.flink.client.program.rest.RestClusterClient.lambda$pollResourceAsync$33(RestClusterClient.java:794)
at java.base/java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:863)
at java.base/java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:841)
at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:510)
at java.base/java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:2179)
at org.apache.flink.util.concurrent.FutureUtils.lambda$retryOperationWithDelay$6(FutureUtils.java:302)
at java.base/java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:863)
at java.base/java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:841)
at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:510)
at java.base/java.util.concurrent.CompletableFuture.postFire(CompletableFuture.java:614)
at java.base/java.util.concurrent.CompletableFuture$UniCompose.tryFire(CompletableFuture.java:1163)
at java.base/java.util.concurrent.CompletableFuture$Completion.run(CompletableFuture.java:482)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642)
at java.base/java.lang.Thread.run(Thread.java:1583)
Caused by: org.apache.flink.runtime.client.JobCancellationException: Job was cancelled.

### Version

3.2.x

### Are you willing to submit PR?

- [ ] Yes I am willing to submit a PR!

### Code of Conduct

- [x] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct)

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.