news 2026/10/12 2:06:08

Apache Beam Python 的 Mean 聚合变换:Globally 与 PerKey 用法及底层实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Python 的 Mean 聚合变换:Globally 与 PerKey 用法及底层实现

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

本文围绕 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.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载
上一篇:RDP Wrapper:如何免费解锁Windows多用户远程桌面限制?
下一篇:RDPWrap完整指南:免费解锁Windows多用户远程桌面的终极解决方案

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/12 2:05:20

开源源测量单元USMU全解析:从硬件拆解到校准实操

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/12 2:03:34

摄影测量三大核心:后方交会、相对定向与光束法平差实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/12 2:02:58

Langchain01_框架之模型的创建与调用

模型创建3种方式 1.使用特定的Model Class(最直接,但不好用) LangChain为一些大模型供应商提供了专门的Model类,导入对应的具体类(如 ChatOpenAI、ChatAnthropic、ChatDeepSeek、ChatOllama、ChatHunyuan、ChatTongy…

作者头像 李华