如何在Apache Beam Python/Dataflow中获取字典列表各列唯一值(非侧输入)
获取Dataflow中字典列表各列唯一值(无需全量侧输入)
你之前用侧输入的方式虽然能实现,但当数据量上来的时候,全量侧输入会把所有数据复制到每个Worker,不仅占用大量内存,还会拖慢整个 pipeline 的执行效率。这里有个更适合分布式场景的方案,用Dataflow的Combine变换来实现,完全不需要全量侧输入:
核心思路
我们可以把问题拆成两步,利用分布式聚合的特性来避免全量数据的传输:
- 拆分键值对:把每个字典展开成
(列名, 值)的元组,这样每一行数据会被拆成多个条目,每个条目对应一列的一个取值。 - 按列分组去重:用
CombinePerKey结合自定义的聚合逻辑,对每个列名下的所有值做分布式去重——每个Worker先处理自己分片内的数据做局部去重,再把局部结果合并成全局唯一值集合。
代码实现(Python示例)
首先定义一个自定义的CombineFn,用来聚合并去重每个列的取值:
import apache_beam as beam from apache_beam.transforms import CombineFn class UniqueValuesCombineFn(CombineFn): # 初始化累加器(用集合来自动去重) def create_accumulator(self): return set() # 把单个值加入累加器 def add_input(self, accumulator, input): accumulator.add(input) return accumulator # 合并多个Worker的局部累加结果 def merge_accumulators(self, accumulators): merged_set = set() for acc in accumulators: merged_set.update(acc) return merged_set # 输出最终的唯一值列表 def extract_output(self, accumulator): return list(accumulator)
然后在主Pipeline里执行流程:
with beam.Pipeline() as p: # 替换成你的实际数据源(比如从BigQuery、文件读取) input_dicts = p | "加载字典列表数据" >> beam.Create([ {"user_id": "u1", "region": "north"}, {"user_id": "u2", "region": "south"}, {"user_id": "u1", "region": "north"}, {"user_id": "u3", "region": "east"} ]) # 把每个字典拆成(列名, 值)元组 key_value_pairs = input_dicts | "展开键值对" >> beam.FlatMap( lambda d: [(col_name, val) for col_name, val in d.items()] ) # 按列名分组,聚合唯一值 unique_per_column = key_value_pairs | "聚合列唯一值" >> beam.CombinePerKey(UniqueValuesCombineFn()) # 输出结果(可以替换成写入存储的逻辑) unique_per_column | "打印结果" >> beam.Map(print)
为什么这个方案更好?
- 分布式执行:CombinePerKey会自动把数据分片到各个Worker,每个Worker只处理自己分片内的数据做局部去重,最后只合并局部的唯一值集合,避免了全量数据的复制和传输。
- 内存友好:不需要在单个Worker里加载所有数据,内存占用随分片大小线性增长,不会出现OOM问题。
- 可扩展:数据量越大,这个方案的性能优势越明显,完全适配Dataflow的分布式架构。
注意事项
- 如果你的值是不可哈希的复杂类型(比如嵌套字典),不能直接用集合去重,需要修改
CombineFn的逻辑,比如用列表存储后再手动去重(或者实现自定义的比较逻辑)。 - 如果需要把结果传递给后续变换使用,可以把聚合后的唯一值集合作为小侧输入(因为聚合后的结果通常远小于原始数据),这样既避免了全量侧输入的问题,又能满足下游的需求。
内容的提问来源于stack exchange,提问作者rachit
相关产品推荐
相关产品推荐

