【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文基于 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) / 2KV("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 从三个层面验证实现:
- 命名正确性:
Mean.globally().getName()为Combine.globally(Mean),Mean.perKey().getName()为Combine.perKey(Mean); - 聚合正确性:通过
CombineFnTester.testCombineFn(Mean.of(), Lists.newArrayList(1, 2, 3, 4), 2.5)验证1,2,3,4的均值为2.5(MeanTest.java); - 编码器正确性:
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 | 通用聚合框架,可传入任意自定义CombineFn | Combine.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.
相关推荐
Apache Beam Java SDK 聚合变换 Mean:全局均值与按键均值的实现与实战
Apache Beam Java SDK 聚合变换 Mean:全局均值与按键均值的实现与实战 Apache Beam 的 Mean 是 Java SDK 中用于
批处理流处理大数据Apache Beam Mean 聚合变换完全指南:全局均值与按 Key 分组均值(Java / Python / Go)
Apache Beam Mean 聚合变换完全指南:全局均值与按 Key 分组均值(Java / Python / Go) 本文以 Apache Beam 官方
批处理流处理大数据Apache Beam Mean 聚合变换:全局均值与按键均值的三语言实战与底层实现解析
Apache Beam Mean 聚合变换:全局均值与按键均值的三语言实战与底层实现解析 本指南聚焦 Apache Beam 内置的 Mean 聚合变换,讲解如
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考