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),仅供参考

相关新闻

Jenkins 2.346.1 内网离线安装插件:依赖解析与版本匹配实战

Jenkins 2.346.1 内网离线安装插件:依赖解析与版本匹配实战

简介:本资源面向在内网、隔离网等无外网环境中部署Jenkins的运维与DevOps工程师,针对Jenkins 2.346.1无法在线拉取插件的问题,提供一套完整的离线插件安装方案。压缩包共约2000个文件,整体314.4MB,涵盖90个jpi与30个hp…

2026/10/12 2:05:08 阅读更多 →
PgDog 配置 crate 的文档注释规范:一份注释同时驱动 Rustdoc 与 JSON Schema

PgDog 配置 crate 的文档注释规范:一份注释同时驱动 Rustdoc 与 JSON Schema

数据库后端 【免费下载链接】pgdog PostgreSQL connection pooler, load balancer and database sharder. 项目地址: https://gitcode.com/gh_mirrors/pg/pgdog 点击查看 免费下载 PgDog 是一个用 Rust 编写的 PostgreSQL 连接池、负载均衡与分库分表中间件&#x…

2026/10/12 2:05:08 阅读更多 →
Learn to Cloud 阶段一:SSH 安全远程接入云端虚拟机——从密钥原理到三大云厂商实操

Learn to Cloud 阶段一:SSH 安全远程接入云端虚拟机——从密钥原理到三大云厂商实操

教程云原生 【免费下载链接】learn-to-cloud A courseware built on the belief that anyone can learn foundational cloud engineering skills with the right guide and discipline 项目地址: https://gitcode.com/gh_mirrors/le/learn-to-cloud 点击查看 免费下…

2026/10/12 2:05:08 阅读更多 →

最新新闻

嵌入式Linux安卓驱动开发:供需、实战与面试全攻略

嵌入式Linux安卓驱动开发:供需、实战与面试全攻略

/* 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:53:39 阅读更多 →
共享Buffer却带宽没降?DDR流量的五大根因与排查实战

共享Buffer却带宽没降?DDR流量的五大根因与排查实战

/* 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:53:39 阅读更多 →
OTFS信道估计实战:压缩感知与相位旋转在高速移动通信中的应用

OTFS信道估计实战:压缩感知与相位旋转在高速移动通信中的应用

/* 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:53:39 阅读更多 →
Qt5.9 C++开发指南章节代码实战:从环境搭建到工程避坑

Qt5.9 C++开发指南章节代码实战:从环境搭建到工程避坑

/* 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:53:39 阅读更多 →
Linux进程虚拟地址空间:从页表映射到段错误排查

Linux进程虚拟地址空间:从页表映射到段错误排查

搞Linux服务端开发的人,迟早会遇到这么一幕:程序跑着跑着突然Segmentation Fault,或者free的时候报double free,又或者top里看到某个进程的VIRT高得离谱,但RES却很低。很多人第一反应是查代码、查日志,但真…

2026/10/12 2:53:39 阅读更多 →
ESP32 上实现 ONVIF 相机:从组件搭建到 NVR 添加实战

ESP32 上实现 ONVIF 相机:从组件搭建到 NVR 添加实战

/* 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:52:39 阅读更多 →

日新闻

复古胶片颗粒感噪点合成器:Canvas ImageData 像素高斯杂色注入算法

复古胶片颗粒感噪点合成器:Canvas ImageData 像素高斯杂色注入算法

在数码相机、高清显示屏与现代矢量图形技术高度发达的今天,画面可以做到绝对的锐利、平滑与无瑕。然而,当一张秋日手账插画或拍立得照片过于“平整无瑕”时,往往会散发出一种冰冷生硬的“数码塑料感(Digital Plasticity&#xff0…

2026/10/12 0:00:59 阅读更多 →
活字印刷古籍线装排版:Canvas 竖排文字与栏线自适应算法

活字印刷古籍线装排版:Canvas 竖排文字与栏线自适应算法

在现代网页与移动端设计中,横排(Horizontal Layout)早已经成为了绝对的主流。然而,当我们翻开泛黄的线装古籍、宋版木刻诗集,或是欣赏一张茶道雅集的手写便签时,那种**自上而下纵向书写、自右向左逐列铺展&…

2026/10/12 0:00:59 阅读更多 →
周日晚间的“精神松绑减震器”:无压力情绪倾倒箱与温和轻声陪伴

周日晚间的“精神松绑减震器”:无压力情绪倾倒箱与温和轻声陪伴

每到周日的晚上八点到十点,很多人心里都会悄悄亮起一盏警示灯。 在心理学上,这种现象有一个专门的称谓——“周日夜晚焦虑症(Sunday Scaries)”。明天又是周一,闹钟又要重新在七点响彻卧房;脑海里仿佛有一个…

2026/10/12 0:00:59 阅读更多 →

周新闻

流感时间序列预测实战:ARIMA/LSTM全流程拆解与避坑指南

流感时间序列预测实战:ARIMA/LSTM全流程拆解与避坑指南

简介:基于 ARIMA、LSTM、Transformer 等模型的流感时间序列预测 Python 源码,面向计算机相关专业课程设计与期末大作业学生,以及项目实战学习者。内容覆盖预处理、平稳性检验、定阶、残差分析、多模型对比预测的完整时序建模流程,…

2026/10/12 0:16:30 阅读更多 →
影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别

影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别

影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别 做影刀RPA自动化,十个新手有八个栽在"往输入框里填东西"这件事上:要么填不进去,要么填了一半,要么直接把原来内容追加在后面。这背后的根因&…

2026/10/12 0:16:38 阅读更多 →
影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容

影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容

影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容 1. 认识影刀:什么场景该用RPA采小说数据 起点中文网的页面结构相对稳定——分类榜单、书籍详情、章节内容三块独立页面,跳转链路清晰。这种场景非常适合影刀自动化&#x…

2026/10/12 0:16:43 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

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

2026/10/11 10:45:37 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

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

2026/10/11 14:36:53 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

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

2026/10/11 14:36:54 阅读更多 →