You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Apache Beam Python/Dataflow中获取字典列表各列唯一值(非侧输入)

获取Dataflow中字典列表各列唯一值(无需全量侧输入)

你之前用侧输入的方式虽然能实现,但当数据量上来的时候,全量侧输入会把所有数据复制到每个Worker,不仅占用大量内存,还会拖慢整个 pipeline 的执行效率。这里有个更适合分布式场景的方案,用Dataflow的Combine变换来实现,完全不需要全量侧输入:

核心思路

我们可以把问题拆成两步,利用分布式聚合的特性来避免全量数据的传输:

  1. 拆分键值对:把每个字典展开成(列名, 值)的元组,这样每一行数据会被拆成多个条目,每个条目对应一列的一个取值。
  2. 按列分组去重:用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 03:52:59