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

Flink从Savepoint恢复时Kafka Sink报UnsupportedVersionException异常求助

环境版本信息
  • CDH版本:6.2.1
  • Flink版本:1.13.1
  • Kafka版本:2.1.0-cdh6.2.1
数据链路

kafka(source) -> flink -> kafka(sink)

问题现象

提交的作业初始运行正常,触发Savepoint后,通过该Savepoint恢复作业时抛出如下异常:

2021-11-22 16:39:52,556 INFO  org.apache.flink.runtime.executiongraph.ExecutionGraph       [] - Job Flink Streaming Job (daf7707813d79b884b2c7b1897801248) switched from state RUNNING to FAILING.
org.apache.flink.runtime.JobException: Recovery is suppressed by FixedDelayRestartBackoffTimeStrategy(maxNumberRestartAttempts=0, backoffTimeMS=3000)
    at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:138) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:82) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:207) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:197) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:188) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:677) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.scheduler.SchedulerNG.updateTaskExecutionState(SchedulerNG.java:79) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:435) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) ~[?:1.8.0_144]
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) ~[?:1.8.0_144]
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[?:1.8.0_144]
    at java.lang.reflect.Method.invoke(Method.java:498) ~[?:1.8.0_144]
    at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:305) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:212) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:77) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:158) ~[data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at scala.PartialFunction$class.applyOrElse(PartialFunction.scala:123) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:170) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.actor.Actor$class.aroundReceive(Actor.scala:517) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.actor.ActorCell.invoke(ActorCell.scala:561) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.Mailbox.run(Mailbox.scala:225) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.Mailbox.exec(Mailbox.scala:235) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
    at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) [data-access-1.0-SNAPSHOT-jar-with-dependencies5.jar:?]
Caused by: org.apache.kafka.common.errors.UnsupportedVersionException: Attempted to write a non-default producerId at version 1

异常核心报错为:尝试在版本1的协议下写入非默认的producerId,版本不支持

根因分析

Flink 1.13版本的Kafka Sink默认开启了幂等生产者特性,该特性依赖Kafka服务端支持producerId相关能力。当前使用的Kafka 2.1.0版本消息格式magic value为1,不支持写入非默认producerId。作业第一次运行时producer是新初始化的没有问题,但Savepoint会持久化Kafka Producer的状态包括已生成的producerId,从Savepoint恢复时会复用该producerId写入Kafka,直接触发版本不兼容报错。

解决方案

可根据实际场景选择以下任意一种方案:

  • 方案1(改动最小,快速修复):在Kafka Sink的配置中关闭幂等性,添加配置项properties.put("enable.idempotence", false)即可
  • 方案2(根本解决):将Kafka集群版本升级至2.2及以上,高版本Kafka的消息格式magic value为2,原生支持producerId相关特性
  • 方案3(保留幂等性临时修复):从Savepoint恢复作业时添加启动参数--allow-non-restored-state,丢弃旧的Kafka Producer状态。注意该操作会导致作业恢复后首次写入可能出现重复数据,需要业务侧自行保证幂等处理。

内容的提问来源于stack exchange,提问作者waynecookie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 03:15:01