Beam中使用Flatten合并PCollection时顺序异常的解决方法咨询
解决PCollection合并后顺序不稳定的问题
Flatten本身是并行处理多个输入集合的操作,不会保证输出顺序,要实现文件内容在前、表数据在后的固定顺序,核心思路是给元素加优先级标记,通过排序来强制顺序,具体实现步骤如下:
给两个PCollection的元素分别添加优先级标记
读取文件内容的集合,给每个元素绑定优先级0(确保排在最前面);读取数据库表数据的集合,给每个元素绑定优先级1。以Beam为例,用Tuple包装元素:// 包装文件内容元素 PCollection<Tuple2<Integer, String>> fileWithPriority = fileContent.apply( MapElements.via((String content) -> Tuple2.of(0, content)) ); // 包装数据库表数据元素 PCollection<Tuple2<Integer, String>> dbWithPriority = dbData.apply( MapElements.via((String row) -> Tuple2.of(1, row)) );合并带优先级的集合
用Flatten合并两个已经标记好优先级的PCollection:PCollection<Tuple2<Integer, String>> merged = PCollectionList .of(fileWithPriority) .and(dbWithPriority) .apply(Flatten.pCollections());按优先级排序并移除标记
通过排序操作让优先级0的元素全部排在1之前,再去掉优先级标记得到最终结果:PCollection<String> finalOutput = merged.apply( Sort.<Tuple2<Integer, String>>byKey() ).apply(MapElements.via((Tuple2<Integer, String> elem) -> elem.f1));
如果是用Flink框架,逻辑完全一致:给元素添加优先级key,通过sortPartition或者全局排序控制顺序,最后移除标记即可。
如果你的场景中文件内容只有单条、表数据量不大,还可以用侧输出流关联的方式:先读取文件内容作为主数据流,将数据库表数据作为侧输出缓存,在处理主元素时直接把所有表数据拼接在后面,这种方式性能更优,但只适合小数据量的表场景。
内容的提问来源于stack exchange,提问作者Priyab
相关产品推荐
相关产品推荐

