【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文围绕 Apache Beam Python SDK 中计算算术平均值的Mean聚合变换展开,讲解如何对整条PCollection使用Mean.Globally()求全局均值、对键值对集合按 key 使用Mean.PerKey()分组求均值,并结合仓库源码剖析其底层的MeanCombineFn、Cython 加速实现与空集合时的NaN语义。读完本文,你将掌握均值聚合的标准写法、可复用的完整示例,以及它在分布式批处理和流处理窗口场景下的行为细节。
Mean 变换是什么
Mean是 Apache Beam Python SDK 中用于计算集合元素算术平均值(arithmetic mean)的聚合变换,定义在 sdks/python/apache_beam/transforms/combiners.py 中。它提供两种用法:
Mean.Globally():计算整个PCollection中所有元素的平均值,输出单个数值;Mean.PerKey():对一个由键值对组成的PCollection,分别计算每个 key 对应的所有 value 的平均值,输出(key, mean)对。
两者的底层都是CombineFn机制:数据先按窗口/键分组,再在分布式执行时通过累加器完成"求和 + 计数"的增量合并,最后一次性输出均值。
示例 1:用Mean.Globally()求整条集合的平均值
Mean.Globally()作用于整条PCollection,返回其中全部元素的算术平均值。下面的示例创建一条管道并输出全局均值:
import apache_beam as beam from apache_beam.transforms import combiners with beam.Pipeline() as pipeline: avg = ( pipeline | 'Create numbers' >> beam.Create([6, 3, 1, 1, 9, 1, 5, 2, 0, 6]) | 'Compute mean' >> combiners.Mean.Globally() | beam.Map(print) )对[6, 3, 1, 1, 9, 1, 5, 2, 0, 6]这 10 个数,输出结果为3.4(总和 34 除以 10)。
这段代码与仓库中的单元测试一一对应:在 sdks/python/apache_beam/transforms/combiners_test.py#L93-L118 的test_builtin_combines中,测试用同样的数据调用combine.Mean.Globally(),并用assert_that(result_mean, equal_to([mean]))验证输出等于sum(vals) / float(len(vals))。
无默认值模式:without_defaults()
在流处理中,一个窗口可能没有收到任何元素。Mean.Globally()默认(has_defaults=True)会对空输入产生一个输出;而调用.without_defaults()后,空集合/空窗口将不产生任何输出。仓库测试 combiners_test.py#L109-L125 展示了典型场景:数据先经WindowInto(FixedWindows(60))分窗,再对每个窗口调用combiners.Mean.Globally().without_defaults()求窗口均值。
import apache_beam as beam from apache_beam.transforms import combiners from apache_beam.transforms import window with beam.Pipeline() as pipeline: windowed_mean = ( pipeline | beam.Create([ window.TimestampedValue(2, 0), window.TimestampedValue(5, 1), window.TimestampedValue(9, 30), ]) | beam.WindowInto(window.FixedWindows(60)) | combiners.Mean.Globally().without_defaults() )对应实现中,Mean.Globally继承自CombinerWithoutDefaults,其expand方法根据has_defaults决定是否在结果上追加without_defaults()语义,见 combiners.py#L90-L98。
示例 2:用Mean.PerKey()按 key 分组求平均值
Mean.PerKey()接收一个键值对PCollection,为每个唯一 key 计算其所有 value 的平均值,输出(key, mean):
import apache_beam as beam from apache_beam.transforms import combiners with beam.Pipeline() as pipeline: mean_per_key = ( pipeline | 'Create key-value pairs' >> beam.Create([ ('a', 1), ('a', 1), ('a', 4), ('b', 1), ('b', 13)]) | 'Compute mean per key' >> combiners.Mean.PerKey() | beam.Map(print) )输出结果:
('a', 2.0) # (1 + 1 + 4) / 3 ('b', 7.0) # (1 + 13) / 2这与仓库测试 combiners_test.py#L603-L620 中的test_MeanCombineFn_combine完全一致:测试构造[('a', 1), ('a', 1), ('a', 4), ('b', 1), ('b', 13)],断言Mean.PerKey()输出[('a', 2), ('b', 7)]。
Mean.PerKey的expand方法内部直接委托给core.CombinePerKey(MeanCombineFn()),见 combiners.py#L100-L103。
底层原理:MeanCombineFn与 Cython 加速
纯 Python 实现:(sum, count)累加器
Mean.Globally()与Mean.PerKey()最终都使用同一个MeanCombineFn(定义于 combiners.py#L110-L134)。它由四个核心方法组成:
| 方法 | 行为 |
|---|---|
create_accumulator() | 初始化累加器(0, 0),即(sum, count) |
add_input(sum_count, element) | 累加sum + element、count + 1,返回新的(sum, count) |
merge_accumulators(accumulators) | 把多个累加器的sum、count分别相加合并 |
extract_output(sum_count) | 若count == 0返回float('NaN'),否则返回sum / float(count) |
这就是分布式聚合的典型三段式:在各 worker 上就地累加、跨 worker 合并累加器、最后提取结果。均值不再需要保存全部元素,而只需维护"总和 + 个数"两个标量,因此内存占用与输入规模无关。
类型分派与 Cython 加速
MeanCombineFn还实现了for_input_type(input_type)(见 combiners.py#L129-L134):当输入类型是int时改用cy_combiners.MeanInt64Fn,是float时改用cy_combiners.MeanFloatFn,否则回退到纯 Python 实现。这些加速版本定义在 sdks/python/apache_beam/transforms/cy_combiners.py:
MeanInt64Accumulator(cy_combiners.py#L164-L193):内部维护整型sum与count,add_input会对元素做int转换并校验是否在INT64_MIN ~ INT64_MAX范围内,越界抛出OverflowError;extract_output在 sum 溢出时先做模2**64回绕再还原符号位,最终用整数除法sum // count得到结果;MeanDoubleAccumulator(cy_combiners.py#L318-L334):浮点版本,add_input把元素转成float后累加,输出时同样在count为 0 时返回NaN;MeanInt64Fn、MeanFloatFn(cy_combiners.py#L253-L364):通过_accumulator_type把上述累加器绑定为AccumulatorCombineFn,让 Cython 编译路径直接操作累加器对象,显著降低逐元素处理的 Python 开销。
空集合的行为:输出NaN
需要特别注意:对空PCollection(或空窗口)求均值时,count == 0,extract_output返回float('NaN')。仓库测试 combiners_test.py#L622-L643 的test_MeanCombineFn_combine_empty专门验证了这一行为:对beam.Create([])求全局均值得到nan(测试用beam.Map(str)把 NaN 转成字符串'nan'再断言,因为 NaN 无法与自身比较),而Mean.PerKey()在空输入下输出空集合。
与 CombineGlobally / CombinePerKey 的关系
Mean是通用合并变换CombineGlobally/CombinePerKey的便捷封装:Mean.Globally()等价于beam.CombineGlobally(MeanCombineFn()),Mean.PerKey()等价于beam.CombinePerKey(MeanCombineFn())。这意味着你也可以直接使用底层的MeanCombineFn与其他CombineFn组合(例如用TupleCombineFn(max, combiners.MeanCombineFn(), sum)在一次扫描中同时求最大值、均值与总和,参见 combiners_test.py#L264-L268),或者利用with_hot_key_fanout/with_fanout对热点 key 与大集合进行扇出优化(见 combiners_test.py#L520-L545)。
相关变换
Mean属于 Apache Beam 的聚合(aggregation)类变换家族,在 Python 文档中与以下变换归为一组:
- CombineGlobally:对整个集合执行任意自定义合并函数;
- CombinePerKey:对键值集合按 key 执行合并;
- Max:求集合最大值;
- Min:求集合最小值;
- Sum:求集合元素之和。
选用建议:当只需要"平均"这一语义时,直接用Mean最简洁;当需要把均值与求和、计数、最值等在一次扫描中一起计算,或需要自定义合并逻辑时,则应改用CombineGlobally/CombinePerKey并传入对应的CombineFn。
小结
Mean.Globally()计算整条PCollection的全局算术平均值,Mean.PerKey()按 key 分组求均值;- 底层统一由
MeanCombineFn实现(sum, count)累加器,并通过for_input_type分派到 Cython 加速的MeanInt64Fn/MeanFloatFn; - 空集合/空窗口的均值输出为
float('NaN'),流处理中可用without_defaults()抑制空窗口输出; - 相关实现与测试可分别查看 combiners.py、cy_combiners.py 与 combiners_test.py,官方 API 参考为
apache_beam.transforms.combiners.Mean。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Java SDK 聚合变换 Mean 深度解析:globally 与 perKey 的用法、原理与源码实现
Apache Beam Java SDK 聚合变换 Mean 深度解析:globally 与 perKey 的用法、原理与源码实现 Apache Beam 提供
大数据批处理流处理数据工程Apache Beam Mean 聚合变换全解析:Globally 与 PerKey 求平均值的跨语言实战指南
Apache Beam Mean 聚合变换全解析:Globally 与 PerKey 求平均值的跨语言实战指南 本文以 Apache Beam 的 Tour o
Apache Beam Python Count 聚合变换详解:Globally / PerKey / PerElement 三种计数方式
Apache Beam Python Count 聚合变换详解:Globally / PerKey / PerElement 三种计数方式 Count 是 Ap
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考