Flink 1.14中TableResult#print与DataStream#print的差异及底层原理问询
Flink 1.14中TableResult#print与DataStream#print的差异机制
1. TableResult#print的底层逻辑
你遇到的3分钟打印一次的现象,核心是Flink SQL流处理的sink输出策略与checkpoint绑定导致的:
- 当通过
tableEnv.executeSql().print()执行无界流SQL查询时,这里的print本质是Flink SQL内置的sink算子。 - 在Flink 1.14中,SQL sink默认开启
table.exec.sink.flush-on-checkpoint: true配置,该配置会让sink缓存的数据仅在checkpoint完成后才批量flush输出。 - 你的checkpoint间隔设置为3分钟,因此每到checkpoint完成时,sink才会把过去3分钟积累的所有10秒一条的数据一次性打印出来,表现为3分钟打印一次。
2. DataStream#print的底层逻辑
而tableEnv.toDataStream(table).print()的即时输出,源于DataStream API的print算子逻辑完全不同:
- DataStream的print算子是无缓存的即时处理算子,每条数据流入算子后,会直接调用标准输出完成打印,不会等待任何外部触发条件。
- 它的执行逻辑和checkpoint完全解耦,source每10秒生成一条数据,就会立即被print算子处理输出,所以表现为10秒打印一次。
额外说明
如果想让TableResult#print也实现类似DataStream的即时输出,可以修改SQL sink的配置:
- 关闭
flush-on-checkpoint:设置table.exec.sink.flush-on-checkpoint: false - 同时配置主动flush间隔:比如
table.exec.sink.flush-interval: 10000(10秒),让sink每隔10秒自动flush一次缓存的数据
内容的提问来源于stack exchange,提问作者gfytd
相关产品推荐
相关产品推荐

