【问题标题】:Flink: Table API copy operators in execution planFlink:执行计划中的表 API 复制运算符
【发布时间】: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


    【解决方案1】:

    在当前状态下,每当您将Table 转换为DataSetDataStream 或将其写入TableSink 时,Table API 都会转换整个查询。在您的程序中,您调用了三次writeToSink,这意味着每次翻译完整的查询。

    但是完整的查询是什么?所有表 API 运算符都已应用于 Table。当您在TableEnvironment 中注册Table 时,它基本上注册为视图,即仅注册其定义(定义表的所有运算符)。因此,当您第二次和第三次调用writeToSink 时,这些运算符会被再次翻译。

    如果您将tableA 转换为DataStream 并在TableEnvironment 中注册DataStream,而不是将其注册为Table,则可以解决此问题。如下所示:

    val tableA = ...
    val streamA = tableA.toDataStream[X] // X should be a case class for rows of tableA
    val tableEnv.registerDataStream("tableA", streamA)
    
    tableEnv.ingest("tableA").writeToSink(sinkTableA) // emit tableA by ingesting the registered DataStream
    

    我知道,这不是很方便,但目前是避免重复翻译表格的唯一方法。

    【讨论】:

    • 当我将 TableA 结果转换为 DataStream 时,我无法使用案例类,因为组和选择字段是我工作中配置文件的可配置性。这就是为什么我使用 Flink Table API 而不是 Flink Streaming 的原因之一,我可以在其中使用案例类来输入和输出数据。所以,我应该使用Row 类进行DataStream 翻译。但是,我不能将Row 用于resul1resul2 的下一个选择(组和选择字段也是可配置的)!
    • 如果我将Rowval resultStream = result.toDataStream[Row]; tableEnv.registerDataStream("TableA", resultStream) 一起用于resul1,则选择出现错误:Exception in thread "main" org.apache.flink.table.api.ValidationException: Cannot resolve [group1] given input [f0, f1, f2, f3, f4]。我尝试将Row 与文件列表tableEnv.registerDataStream("TableA", resultStream, fields) 一起使用,但出现错误Exception in thread "main" org.apache.flink.table.api.TableException: Source of type Row(f0: Timestamp, f1: String, f2: String, f3: Long, f4: Long) cannot be converted into Table.
    猜你喜欢
    • 2018-04-10
    • 1970-01-01
    • 2021-05-17
    • 2016-07-12
    • 2015-05-28
    • 1970-01-01
    • 1970-01-01
    • 2023-04-08
    • 2013-07-23
    相关资源
    最近更新 更多