Flink 1.20从Savepoint恢复时Kafka Sink任务报EndOfInput异常
Flink SQL从Savepoint恢复时Kafka Sink抛出Received element after endOfInput异常
环境信息
- Flink版本:1.20.1
- Flink Kafka连接器版本:3.3.0
- JDK版本:1.8
对应的Flink SQL任务
CREATE TEMPORARY TABLE `t1` ( `org_id` VARCHAR ) WITH ( 'connector' = 'datagen', 'rows-per-second' = '1' ); CREATE TEMPORARY TABLE `kafka_out` ( `org_id` VARCHAR ) WITH ( 'connector' = 'kafka', 'properties.bootstrap.servers' = 'xxx:9092', 'topic' = 'test2', 'value.format' = 'json', 'properties.enable.idempotence' = 'false' ); insert into kafka_out select * from t1;
复现步骤
- 启动该SQL任务,通过Savepoint停止任务
- 使用该Savepoint重启任务,立即抛出以下异常:
2025-12-08 17:14:44.002 [kafka_out[2]: Writer -> kafka_out[2]: Committer (1/1)#0] WARN org.apache.flink.runtime.taskmanager.Task - kafka_out[2]: Writer -> kafka_out[2]: Committer (1/1)#0 (54d08e14b76ac2f5d7f195c7732acee6_20ba6b65f97481d5570070de90e4e791_0_0) switched from RUNNING to FAILED with failure cause: java.lang.IllegalStateException: Received element after endOfInput: Record @ (undef) : org.apache.flink.table.data.binary.BinaryRowData@3f797ef9 at org.apache.flink.util.Preconditions.checkState(Preconditions.java:215) at org.apache.flink.streaming.runtime.operators.sink.SinkWriterOperator.processElement(SinkWriterOperator.java:206) at org.apache.flink.streaming.runtime.tasks.OneInputStreamTask$StreamTaskNetworkOutput.emitRecord(OneInputStreamTask.java:238) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.processElement(AbstractStreamTaskNetworkInput.java:157) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.emitNext(AbstractStreamTaskNetworkInput.java:114) at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:638) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:973) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:917) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:970) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:949) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:763) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575) at java.lang.Thread.run(Thread.java:748)
问题原因分析
通过分析Flink 1.20的SinkWriterOperator源码,发现异常根源如下:
- 触发Savepoint时,SourceOperator调用stop方法,通过
emitNextNotReading发送DataInputStatus.END_OF_DATA状态 - 该状态最终触发
SinkWriterOperator.endInput方法,将endOfInput变量设为true,调用栈:
org.apache.flink.streaming.runtime.operators.sink.SinkWriterOperator.endInput(SinkWriterOperator.java:232) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.endOperatorInput(StreamOperatorWrapper.java:101) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$finish$0(StreamOperatorWrapper.java:154) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.finish(StreamOperatorWrapper.java:154) at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.finish(StreamOperatorWrapper.java:161) at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.finishOperators(RegularOperatorChain.java:115) at org.apache.flink.streaming.runtime.tasks.StreamTask.endData(StreamTask.java:695) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:653) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:973) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:917) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:970) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:949) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:763) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575) at java.lang.Thread.run(Thread.java:748)
endOfInput变量会被持久化到Savepoint的状态中- 从Savepoint恢复时,该变量被加载为true,当
SinkWriterOperator处理新输入记录时,会触发检查抛出异常,对应源码:
@Override public void processElement(StreamRecord<InputT> element) throws Exception { checkState(!endOfInput, "Received element after endOfInput: %s", element); context.element = element; sinkWriter.write(element.getValue(), context); }
解决方法
方法一:修改Savepoint状态(临时应急)
使用Flink的State Processor API修改Savepoint中的endOfInput变量值为false:
- 编写离线程序加载目标Savepoint
- 定位到
SinkWriterOperator对应的状态,将endOfInput变量修改为false - 重新生成修改后的Savepoint,用它恢复任务
方法二:升级Flink版本
该问题是Flink 1.20.x版本的已知缺陷,后续版本(如1.21.0及以上)已修复了Savepoint中持久化endOfInput变量的逻辑,升级到最新稳定版可彻底解决问题。
方法三:调整Savepoint触发方式
触发Savepoint时使用--drain参数:
flink savepoint <job-id> <savepoint-path> --drain
该参数会让任务处理完所有缓存数据后再停止,避免触发END_OF_DATA状态,适合需要确保数据不丢失的场景,但任务停止时间可能更长。
内容的提问来源于stack exchange,提问作者user32053573
相关产品推荐
相关产品推荐

