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)
根因分析
该异常由版本兼容性问题导致:
- Upsert Kafka Sink的buffer flush缓冲特性是从Flink 1.14版本才正式支持的,低于1.14的Flink版本没有实现对应特性的新Sink接口
- 当低版本Flink配置了buffer-flush相关参数时,框架会尝试加载新版Sink的Transformation,而低版本仅存在Legacy Sink实现,因此触发找不到对应Transformation的异常
- 若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
相关产品推荐
相关产品推荐

