Flink1.13如何将DataSet API替换为DataStream API实现first取数打印
Flink 1.13 DataStream 实现全局取指定条数打印
原DataSet API的first(limit)是全局维度截取前limit条数据(未提前排序时不保证返回数据的顺序,凑够条数即返回),DataStream API没有直接提供同名算子,按照下面的写法可以实现完全等价的效果:
实现代码
// 替换YourDataType为你实际的数据流数据类型,limit为要截取的条数 dataStream // 将所有上游数据全局路由到下游同一个并行实例 .global() .flatMap(new RichFlatMapFunction<YourDataType, YourDataType>() { // 计数器,统计已经输出的条数 private int currentCount = 0; @Override public void flatMap(YourDataType record, Collector<YourDataType> out) { if (currentCount < limit) { out.collect(record); currentCount++; } // 超过条数后直接丢弃后续数据,不做任何处理 } }) // 强制该截断算子并行度为1,保证全局计数准确 .setParallelism(1) .print();
注意事项
- 必须搭配
global()数据路由和setParallelism(1)使用:如果保留多并行度,每个并行子任务会独立计数,最终输出的总条数会变成并行度 * limit,不符合预期。 - 该实现和原DataSet的
first(limit)行为完全对齐:不需要提前排序,只要凑够指定条数就输出,后续数据全部丢弃。如果需要取排序后的前N条,可以在调用global()之前先完成全局排序逻辑,再走后面的截断即可。 - 全局路由到单并行度的算子会产生数据倾斜,这是全局取前N条的固有开销,原DataSet的
first算子底层也是通过单节点汇总截断实现的,不存在性能更好的无倾斜实现方式。
如果你的需求是按键分组后取每个Key对应的前N条,只需要去掉global(),在keyBy(你的key选择器)返回的KeyedStream上写同样的计数截断逻辑即可,不需要强制设并行度为1,Flink会保证同一个Key的数据路由到同一个task,计数准确。
内容的提问来源于stack exchange,提问作者xss
相关产品推荐
相关产品推荐

