Flink 1.14两类Row数据流转Table字段别名异常原因咨询
Flink 1.14中DataStream转Table的字段别名差异原因
一、fromElements创建的DataStream能正常执行的原因
当使用env.fromElements()创建DataStream时,Flink可以在编译阶段直接分析传入的Row实例(比如Row.of("Alice", 12)),自动生成包含字段数量、字段类型的TypeInformation<Row>。此时调用.as("name", "score")时,指定的字段数量(2个)和Flink推断出的Row字段数量完全匹配,所以能正常执行。
对应的代码片段:
DataStream<Row> dataStream = env.fromElements( Row.of("Alice", 12), Row.of("Bob", 10), Row.of("Alice", 100)); Table table = tableEnv.fromDataStream(dataStream).as("name", "score");
二、自定义Source的DataStream抛出错误的原因
自定义SourceFunction生成的DataStream,Flink在编译阶段无法提前获取Row的结构信息——因为Source的数据是运行时才产生的,编译期不知道Row包含多少字段、每个字段的类型是什么。这种情况下,Flink会默认将Row的类型标记为GenericType<Row>,而GenericType不会记录Row的字段数量等元数据。当调用.as("name", "score")指定2个字段别名时,Flink无法确认Row是否有足够的字段,因此抛出错误Aliasing more fields than we actually have。
对应的代码片段:
DataStreamSource<Row> ds = env.addSource(new MockSource()); Table table = tableEnv.fromDataStream(ds).as("name", "score"); ...... public static class MockSource implements SourceFunction<Row> { private volatile boolean running = true; @Override public void run(SourceContext<Row> ctx) throws Exception { while (running) { Row row = Row.of("Alice", 12); ctx.collect(row); Thread.sleep(10000); } } @Override public void cancel() { running = false; } }
补充:解决自定义Source的类型推断问题
如果要让自定义Source的DataStream正常转Table,需要显式指定Row的TypeInformation:
// 显式定义Row的类型信息:2个字段,分别为String类型和Integer类型 TypeInformation<Row> rowType = Types.ROW(Types.STRING(), Types.INT()); DataStreamSource<Row> ds = env.addSource(new MockSource()).returns(rowType); Table table = tableEnv.fromDataStream(ds).as("name", "score");
内容的提问来源于stack exchange,提问作者gfytd
相关产品推荐
相关产品推荐

