Flink有界流输出乱序原因、有序配置及多Sink疑问
问题解答
1. 有界流输出顺序混乱的原因及有序输出配置
输出顺序混乱是因为Flink默认以并行方式处理数据:
你的Source并行度是1,但print()算子默认使用全局并行度(本地模式下通常等于机器CPU核心数),数据会被分区发送到不同的Sink并行子任务处理。每个子任务独立输出结果,各子任务的输出交织在一起,就会出现顺序混乱的情况。
要实现有序输出,有两种简单方法:
- 方法一:设置
print()算子的并行度为1,强制串行输出:dataStream.print().setParallelism(1); - 方法二:设置全局执行环境的并行度为1,让整个作业所有算子都串行执行:
env.setParallelism(1);
2. 日志显示8个Sink的原因
print()算子作为Sink,它的并行度默认继承Flink的全局默认并行度。在本地运行模式下,全局默认并行度等于你的机器CPU核心数(这里是8核),所以Flink会启动8个Sink并行实例,对应日志里的Sink: Unnamed(1/8)到Sink: Unnamed(8/8)。
如果想修改Sink的数量,直接调整print()算子的并行度即可,比如设置为2:
dataStream.print().setParallelism(2);
内容的提问来源于stack exchange,提问作者overexchange
相关产品推荐
相关产品推荐

