Apache Beam Python 的 Mean 聚合变换:Globally 与 PerKey 用法及底层实现
【免费下载链接】beamApache 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_defaultsTrue会对空输入产生一个输出而调用.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.pyMeanInt64Accumulatorcy_combiners.py#L164-L193内部维护整型sum与countadd_input会对元素做int转换并校验是否在INT64_MIN ~ INT64_MAX范围内越界抛出OverflowErrorextract_output在 sum 溢出时先做模2**64回绕再还原符号位最终用整数除法sum // count得到结果MeanDoubleAccumulatorcy_combiners.py#L318-L334浮点版本add_input把元素转成float后累加输出时同样在count为 0 时返回NaNMeanInt64Fn、MeanFloatFncy_combiners.py#L253-L364通过_accumulator_type把上述累加器绑定为AccumulatorCombineFn让 Cython 编译路径直接操作累加器对象显著降低逐元素处理的 Python 开销。空集合的行为输出NaN需要特别注意对空PCollection或空窗口求均值时count 0extract_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。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐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 oApache Beam Python Count 聚合变换详解Globally / PerKey / PerElement 三种计数方式Apache Beam Python Count 聚合变换详解Globally / PerKey / PerElement 三种计数方式 Count 是 Ap上一篇RDP Wrapper如何免费解锁Windows多用户远程桌面限制下一篇RDPWrap完整指南免费解锁Windows多用户远程桌面的终极解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考