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

Flink upsert-kafka sink配置buffer-flush报LegacySinkTransformation不存在错误

异常现象

运行Flink任务时,为Upsert Kafka Sink配置buffer-flush相关参数后启动抛出异常,相同任务逻辑不配置buffer-flush参数时可正常运行。

异常堆栈
Exception in thread "main" java.lang.IllegalStateException: There is no the LegacySinkTransformation.
at org.apache.flink.streaming.api.datastream.DataStreamSink.getTransformation(DataStreamSink.java:71)
at org.apache.flink.table.planner.plan.nodes.exec.common.CommonExecSink.applySinkProvider(CommonExecSink.java:294)
at org.apache.flink.table.planner.plan.nodes.exec.common.CommonExecSink.createSinkTransformation(CommonExecSink.java:145)
at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecSink.translateToPlanInternal(StreamExecSink.java:140)
at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:134)
at org.apache.flink.table.planner.delegation.StreamPlanner.$anonfun$translateToPlan$1(StreamPlanner.scala:71)
at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:233)
at scala.collection.Iterator.foreach(Iterator.scala:937)
at scala.collection.Iterator.foreach$(Iterator.scala:937)
at scala.collection.AbstractIterator.foreach(Iterator.scala:1425)
at scala.collection.IterableLike.foreach(IterableLike.scala:70)
at scala.collection.IterableLike.foreach$(IterableLike.scala:69)
at scala.collection.AbstractIterable.foreach(Iterable.scala:54)
at scala.collection.TraversableLike.map(TraversableLike.scala:233)
at scala.collection.TraversableLike.map$(TraversableLike.scala:226)
at scala.collection.AbstractTraversable.map(Traversable.scala:104)
at org.apache.flink.table.planner.delegation.StreamPlanner.translateToPlan(StreamPlanner.scala:70)
at org.apache.flink.table.planner.delegation.PlannerBase.translate(PlannerBase.scala:185)
at org.apache.flink.table.api.internal.TableEnvironmentImpl.translate(TableEnvironmentImpl.java:1665)
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:752)
at org.apache.flink.table.api.internal.StatementSetImpl.execute(StatementSetImpl.java:124)
at com.company.flink.FlinkJob.main(FlinkJob.java:260)
根因分析

该异常由版本兼容性问题导致:

  1. Upsert Kafka Sink的buffer flush缓冲特性是从Flink 1.14版本才正式支持的,低于1.14的Flink版本没有实现对应特性的新Sink接口
  2. 当低版本Flink配置了buffer-flush相关参数时,框架会尝试加载新版Sink的Transformation,而低版本仅存在Legacy Sink实现,因此触发找不到对应Transformation的异常
  3. 若Flink核心版本和Upsert Kafka连接器版本不一致,也会触发该类兼容异常
解决方案
  • 方案一:升级全组件版本到Flink 1.14及以上,确保Flink核心、Table Planner、Upsert Kafka连接器的版本完全统一,即可正常使用buffer-flush配置
  • 方案二:若暂时无法升级Flink版本,直接移除所有buffer-flush相关配置参数,回退到无缓冲提交模式即可正常运行任务
buffer-flush配置说明
  • sink.buffer-flush.max-rows:Sink端缓冲的最大数据行数,达到该阈值后会异步将缓冲数据写入Kafka,默认值为0,代表关闭缓冲功能
  • sink.buffer-flush.interval:缓冲的定时刷写间隔,即使缓冲行数未达到阈值,到达间隔时间也会自动刷写数据,默认值为0,代表关闭定时刷写
  • 两个参数只要任意一个配置为大于0的值,缓冲功能就会开启;两个参数均设为0时完全关闭缓冲

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 08:54:03