【问题标题】:Error: There is no the LegacySinkTransformation Flink错误:没有 LegacySinkTransformation Flink
【发布时间】:2021-11-08 04:45:26
【问题描述】:

当使用带有buffer-flush 选项的upsert-kafka sink 时,我遇到了错误,而没有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)

https://ci.apache.org/projects/flink/flink-docs-master/docs/connectors/table/upsert-kafka/#sink-buffer-flush-max-rows

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    这是 Flink 1.14 中的一个严重错误,可能会在下一个补丁版本中修复。更多信息可以在这里找到:

    https://issues.apache.org/jira/browse/FLINK-24596

    一种解决方法是使用最新的 1.13 版本或暂时禁用缓冲。

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-05-06
    • 1970-01-01
    • 2018-05-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多