【发布时间】:2017-09-18 18:02:28
【问题描述】:
我正在尝试在带有 Java 1.8.0_144 IDE Eclipse Mars 的 Windows 10 上使用 Apache Flink 1.3.2 实现日志分析器。
上下文:
- LogMessage 有多种类型。
- 为每种类型创建 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