You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink 1.20从Savepoint恢复时Kafka Sink任务报EndOfInput异常

环境信息

  • Flink版本:1.20.1
  • Flink Kafka连接器版本:3.3.0
  • JDK版本:1.8
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;

复现步骤

  1. 启动该SQL任务,通过Savepoint停止任务
  2. 使用该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源码,发现异常根源如下:

  1. 触发Savepoint时,SourceOperator调用stop方法,通过emitNextNotReading发送DataInputStatus.END_OF_DATA状态
  2. 该状态最终触发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)
  1. endOfInput变量会被持久化到Savepoint的状态中
  2. 从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:

  1. 编写离线程序加载目标Savepoint
  2. 定位到SinkWriterOperator对应的状态,将endOfInput变量修改为false
  3. 重新生成修改后的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.11 21:13:11