【发布时间】:2017-05-27 05:53:55
【问题描述】:
我使用 Flink 1.2.0 Table API 来处理一些流数据。以下是我的代码:
val dataTable = myDataStream
// table A
val tableA = dataTable
.window(Tumble over 5.minutes on 'rowtime as 'w)
.groupBy("w, group1, group2")
.select("w.start as time, group1, group2, data1.sum as data1, data2.sum as data2")
tableEnv.registerTable("tableA", tableA)
// table A sink
tableA.writeToSink(sinkTableA)
//...
// I shoul get some other outputs from TableA output
//...
val dataTable = tableEnv.ingest("tableA")
// table result1
val result1 = dataTable
.window(Tumble over 5.minutes on 'rowtime as 'w)
.groupBy("w, group1")
.select("w.start as time, group1, data1.sum as data1")
// result1 sink
result2.writeToSink(sinkResult1)
// table result2
val result2 = dataTable
.window(Tumble over 5.minutes on 'rowtime as 'w)
.groupBy("w, group2")
.select("w.start as time, group2, data2.sum as data1")
// result2 sink
result2.writeToSink(sinkResult2)
我等待在 flink 执行计划中获取这棵树。 与我在其他 Flink 工作中的 Flink Streaming 相同。
DataStream_Operators -> TableA_Operators -> TableA_Sink
|-> Result1_Operators -> Result1_Sink
|-> Result2_Operators -> Result2_Sink
但是,我用 TableA 的 3 个相同操作符得到了这个!
DataStream_Operators -> TableA_Operators -> TableA_Sink
|-> Copy_of_TableA_Operators -> Result1_Operators -> Result1_Sink
|-> Copy_of_TableA_Operators -> Result2_Operators -> Result2_Sink
我在这项工作的大量输入数据中表现不佳。
我该如何解决这个问题并获得最佳执行计划?
我不明白,Flink Table API 和 SQL 是什么实验性功能和 也许它会在下一个版本中修复。
【问题讨论】:
标签: scala apache-flink