【问题标题】:is it not allowed to query on supertype of POJO Dataset in Apache Flink Table API是否不允许在 Apache Flink Table API 中查询 POJO 数据集的超类型
【发布时间】:2017-09-18 18:02:28
【问题描述】:

我正在尝试在带有 Java 1.8.0_144 IDE Eclipse Mars 的 Windows 10 上使用 Apache Flink 1.3.2 实现日志分析器。

上下文:

  • Lo​​gMessage 有多种类型。
  • 为每种类型创建 POJO。
  • 为每种类型创建 POJO 类型的 DataSet 实例。
  • 然后使用表API查询如下。

这很好用。

DataSet<String> rawLogs = env.readTextFile(input);// input is the data file path
DataSet<FirstBackupMessage> logMsgPOJODataSet = rawLogs.map(new LogMapFunction());
BatchTableEnvironment tableEnv = TableEnvironment.getTableEnvironment(env); 
Table LogMessageTable = tableEnv.fromDataSet(logMsgPOJODataSet);
Table result = tableEnv .sql("Select taskId from " + LogMessageTable);
tableEnv.toDataSet(result, Row.class).print();

要求: 我正在尝试使用工厂模型来概括此实现。 为了做到这一点,我尝试将 POJO 类推广到 LogMessage interface。在上述情况下:

public class FirstBackupMessage implements LogMessage
similarly 
public class SecondBackupMessage implements LogMessage
public class ThirdBackupMessage implements LogMessage

在 MapFunction 实现中,我正在填充特定的类实例,但 map 函数的输出映射到通用引用,即 LogMessage 在上述情况下,它将是

DataSet<LogMessage> logMsgPOJODataSet = rawLogs.map(new LogMapFunction());  
//the LogMapFunction.map method is populating FirstBackupMessage

在此之后,如果我尝试查询 POJO FirstBackupMessage 中存在的字段,但现在参考接口,即 LogMessage 它抛出异常说我正在查询的字段未找到。

但是

奇怪的是,如果我打印带有通用引用的 DataSet,即 logMsgPOJODataSet.print(),它会打印特定 POJO 中的所有字段,在这种情况下为 FirstBackupMessage。

问题: 在 Flink Table API 中是否允许 / 使用这种将泛型引用转换为 DataSet 的类型?

【问题讨论】:

  • 这是OO的核心!在 Java 中,变量有一个类型,当取消引用这个变量时,您只能访问该声明的类型的成员。另一方面,在运行时此变量可以引用任何子类型的实例,并且方法调用(即toString)由此运行时类型确定。
  • 感谢您的回复!是的,你是对的,我的错。这回答了为什么 toString 打印所有值的第二部分。但接下来的问题将是(或应该是)如何在运行时在 flink 数据集或 Dayastream API 中将其转换回子类型。

标签: java apache-flink


【解决方案1】:

Table API/SQL 库对关系表进行操作。通过调用TableEnvironment.fromDataSet(logMsgPOJODataSet)DataSetlogMsgPOJODataSet被逻辑转换成一个表。在这个过程中,需要根据logMsgPOJODataSetDataSet的类型来识别新表的schema。 Flink 的 DataSet API 使用TypeInformation 来判断DataSet 的数据类型。

由于logMsgPOJODataSet DataSet 的类型是LogMessage,Table API 不知道它的任何子类型。因此,LogMessage 的所有字段都包括在内,但没有子类型字段。

无论如何,在同一个表中处理不同类型的行是不可能的。所有行必须具有相同的架构。处理这种情况的两种方法是:

  1. 使架构成为所有子类型的超集,并为不受支持的类型设置空值。可能会添加另一个指示子类型的字段。
  2. 添加一个包含所有子类型数据的通用 Map&lt;String, String&gt; 字段。

在这两种情况下,都需要使用 DataSet API 完成转换,例如使用 MapFunction

【讨论】:

    猜你喜欢
    • 2018-05-28
    • 2022-10-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-18
    • 1970-01-01
    相关资源
    最近更新 更多