☰
Apache Beam Java Mean 变换详解:全局均值与按键均值聚合
2026/10/12 3:51:58 网站建设 项目流程

【免费下载链接】beam

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

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

本文基于 Apache Beam 仓库内 Java SDK 文档与核心源码,系统讲解聚合变换Mean的两种用法:Mean.globally()(对整条PCollection计算算术平均值)与Mean.perKey()(对KV集合按 Key 分组求均值)。读者将掌握调用方式、空输入行为、底层CombineFn累加器与编码器原理,并获得可直接运行的示例代码。

一、Mean是什么

Mean是 Apache Beam Java SDK 提供的一组PTransform,用于计算集合中元素的算术平均值(arithmetic mean,即平均数):

  • Mean.globally():对输入PCollection<NumT>的全部元素求均值,返回只含一个元素的PCollection<Double>。
  • Mean.perKey():对输入PCollection<KV<K, NumT>>按不同 Key 分组,输出PCollection<KV<K, Double>>,每个 Key 对应其关联值的均值。

从源码结构看,Mean本身只是一个命名空间(namespace)类,其构造函数为私有、不可实例化,全部能力通过两个静态工厂方法暴露,定义在 Mean.java:

  • Mean.globally()等价于Combine.globally(Mean.of()),返回Combine.Globally<NumT, Double>;
  • Mean.perKey()等价于Combine.perKey(Mean.of()),返回Combine.PerKey<K, NumT, Double>。

也就是说,Mean是基于Combine构建的"现成聚合函数",其通用性由Combine提供(支持分布式部分聚合、窗口化等),而"求均值"这一具体运算规则则由内部MeanFn(一个Combine.AccumulatingCombineFn)实现。这一点可以通过单元测试得到印证:MeanTest.java 断言Mean.globally().getName()为"Combine.globally(Mean)"、Mean.perKey().getName()为"Combine.perKey(Mean)"。

二、Mean.globally():全局均值

2.1 用法与签名

PCollection<Long> input = ...; PCollection<Double> mean = input.apply(Mean.<Long>globally());

方法签名为:

public static <NumT extends Number> Combine.Globally<NumT, Double> globally()

要点:

  • 输入元素类型NumT必须是java.lang.Number的子类型(Integer、Long、Double、Float、BigDecimal等均可);
  • 输出恒为PCollection<Double>,即均值以Double表示;
  • 全局聚合的结果是单元素集合——一条输入PCollection对应一条包含均值的输出元素。

2.2 空输入行为

官方文档与globally()的 Javadoc 均声明:若输入集合中没有元素,则返回0(见 mean.md 与 Mean.java)。

不过需要留意当前仓库源码中的一个细节:Mean.of()返回的MeanFn在累加器计数为 0 时,其extractOutput()实际返回的是Double.NaN(Mean.java),且Mean.of()的 Javadoc 明确写着 "ReturnsDouble.NaNif combining zero elements"。二者看似矛盾,实则是两条不同代码路径:

  • 直接使用Mean.of()作为CombineFn且输入为空时,聚合结果由extractOutput()决定,即NaN;
  • 使用Mean.globally()时,Combine.Globally默认开启insertDefault,会在输出为空时插入默认值——该默认值来自fn.defaultValue(),即对空累加器调用extractOutput()的结果(Combine.java)。

因此在编写依赖空集合行为的数据流时,建议以当前所用 Beam 版本的实际运行结果为准,并可通过Mean.globally().withoutDefaults()关闭默认值注入(该选项在输入非全局窗口化、且输出不作为侧输入(side input)时是必须的,见 Combine.java)。

2.3 完整可运行示例

仓库 examples/java 中提供了可直接运行的 MeanExample.java:

PipelineOptions options = PipelineOptionsFactory.create(); Pipeline pipeline = Pipeline.create(options); PCollection<Double> pc = pipeline.apply(Create.of(1.0, 2.0, 3.0, 4.0, 5.0)); PCollection<Double> mean = pc.apply(Mean.globally()); // Log values mean.apply(ParDo.of(new LogOutput<>("PCollection numbers after Mean transform: "))); pipeline.run();

输入1.0, 2.0, 3.0, 4.0, 5.0,输出为3.0。该示例同时以 Playground 元数据标注(name: Mean、complexity: BASIC、categories: Core Transforms),与文档页中的SDK_JAVA_Mean在线示例一一对应。

三、Mean.perKey():按键求均值

3.1 用法与签名

PCollection<KV<String, Integer>> input = ...; PCollection<KV<String, Double>> meanPerKey = input.apply(Mean.<String, Integer>perKey());

方法签名为:

public static <K, NumT extends Number> Combine.PerKey<K, NumT, Double> perKey()

要点:

  • 输入为PCollection<KV<K, NumT>>,其中K为任意可编码的 Key 类型,NumT仍须是Number子类型;
  • 输出为PCollection<KV<K, Double>>,每个不同 Key 对应一个元素,值为该 Key 下所有输入值(经doubleValue()转换后)的均值;
  • 其语义与GroupByKey后对每组Iterable手工求均值等价,但底层走Combine的执行路径,可实现分布式部分聚合,避免把所有同 Key 数据集中到单台机器处理(这正是 Combine 文档 强调的通信开销优势)。

按时间戳、窗口(bucketing)等语义,perKey()的行为与Combine.PerKey完全一致,可参阅 Mean.java 中对Combine.PerKey的引用说明。

3.2 完整可运行示例

仓库中对应的可运行示例为 MeanPerKeyExample.java:

PipelineOptions options = PipelineOptionsFactory.create(); Pipeline pipeline = Pipeline.create(options); PCollection<KV<String, Integer>> input = pipeline.apply( Create.of(KV.of("a", 1), KV.of("a", 2), KV.of("b", 3), KV.of("b", 4), KV.of("b", 5))); PCollection<KV<String, Double>> meanPerKey = input.apply(Mean.perKey()); // Log values meanPerKey.apply(ParDo.of(new LogOutput<>("PCollection numbers after Mean transform: "))); pipeline.run();

输入数据:Keya关联1, 2,Keyb关联3, 4, 5。预期输出:

  • KV("a", 1.5):(1 + 2) / 2
  • KV("b", 4.0):(3 + 4 + 5) / 3

四、源码探秘:均值是怎么算出来的

Mean的底层核心是MeanFn与累加器CountSum,全部定义在 Mean.java 中。

4.1 累加器:只记"个数 + 和"

为避免把所有原始元素传到最终聚合点,MeanFn采用Combine.AccumulatingCombineFn,累加器CountSum只维护两个字段(Mean.java):

long count = 0; // 元素个数 double sum = 0.0; // 元素之和

累加器提供的三个核心方法定义了完整的分布式聚合协议:

  • addInput(NumT element):计数加 1,sum += element.doubleValue()——所有数值类型先统一转为double参与运算(Mean.java);
  • mergeAccumulator(CountSum accumulator):把另一个累加器的count与sum分别相加——这保证了合并的结合律(associativity),是Combine能在多台机器上并行部分聚合的理论前提(Mean.java);
  • extractOutput():count == 0 ? Double.NaN : sum / count,即最终一步才做除法、得到均值(Mean.java)。

4.2 累加器的序列化编码

由于累加器需要跨 worker 传输(合并阶段可能发生在不同于输入阶段所在的机器),MeanFn通过getAccumulatorCoder返回专用的CountSumCoder。该编码器继承AtomicCoder,内部使用BigEndianLongCoder编码count、DoubleCoder编码sum,并实现了verifyDeterministic()以验证编码的确定性(Mean.java)——确定性是 Beam 分布式执行正确性(如精确一次的 shuffle、可复现结果)的重要保障。

4.3 测试验证

MeanTest.java 从三个层面验证实现:

  1. 命名正确性:Mean.globally().getName()为Combine.globally(Mean),Mean.perKey().getName()为Combine.perKey(Mean);
  2. 聚合正确性:通过CombineFnTester.testCombineFn(Mean.of(), Lists.newArrayList(1, 2, 3, 4), 2.5)验证1,2,3,4的均值为2.5(MeanTest.java);
  3. 编码器正确性:coderDecodeEncodeEqual验证编码后再解码与原始值相等,coderSerializable验证编码器可序列化(MeanTest.java)。

五、Mean与Combine的关系

Mean在本质上是Combine的一个"预置特例",二者关系可以从两个层面理解:

类型层面:Mean.globally()的返回类型就是Combine.Globally,Mean.perKey()的返回类型就是Combine.PerKey,因此Combine提供的下列能力对Mean同样可用:

  • withFanout(int fanout):在最终全局合并前增加中间节点,把数据切分到fanout个中间 Key 上做部分聚合,以降低全局聚合步骤的单点负载(Combine.java);
  • asSingletonView():将全局聚合结果作为单元素侧输入(PCollectionView)供其他变换读取(Combine.java);
  • withoutDefaults():关闭空输入时的默认值插入(Combine.java)。

语义层面:Mean的运算规则(CountSum累加器)是"可结合、可交换"的——先分组加和、再合并加和、最后统一相除,结果与顺序无关。这正是Combine要求 CombineFn 满足的代数性质,也是它比GroupByKey + ParDo手工求均值性能更优的原因:可以提前做部分聚合,大幅减少数据洗牌(shuffle)量。详细对比可参见 Combine 变换文档。

六、相关变换一览

Mean是 Beam 聚合变换家族的一员,同一目录下还有语义相近的现成聚合,建议按需选用:

变换作用仓库实现
Max求集合/每组最大值Max.java
Min求集合/每组最小值Min.java
Sum求集合/每组之和Sum.java
Combine通用聚合框架,可传入任意自定义CombineFnCombine.java

对应的官方文档分别位于 max.md、min.md、combine.md。选型建议:

  • 只需要最大值、最小值、求和、均值这类内置运算,直接用Max、Min、Sum、Mean最省事;
  • 需要自定义聚合逻辑(如加权平均、方差、去重计数等),则通过Combine.globally(...)/Combine.perKey(...)传入自定义CombineFn;
  • 需要按 Key 聚合时,优先考虑上述变换的perKey()变体,而不是GroupByKey后手工循环,后者会把同 Key 全部数据集中到单个 worker 处理,带来不必要的通信开销。

七、总结

Mean是 Apache Beam Java SDK 中"开箱即用"的均值聚合变换:

  • Mean.globally()求整条PCollection的算术平均,输出单元素PCollection<Double>;
  • Mean.perKey()按 Key 分组求平均,输出PCollection<KV<K, Double>>;
  • 底层由MeanFn + CountSum累加器实现"只传个数与和、最后相除"的分布式安全算法,并配有专用确定性编码器CountSumCoder;
  • 空输入时,Mean.of()的累加器返回Double.NaN,而Mean.globally()默认会通过Combine的默认值机制插入默认输出,实际行为以所用版本为准;
  • 所有与Combine相关的能力(withFanout、asSingletonView、withoutDefaults等)对Mean同样适用。

若需查看更完整的可运行代码,可研读 MeanExample.java 与 MeanPerKeyExample.java,并结合 MeanTest.java 理解其边界行为。

【免费下载链接】beam

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

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

相关推荐

上一篇:Metallb国际化支持:文档翻译与多语言适配指南
下一篇:tuigreet完全指南:打造高效控制台登录体验的终极方案

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询