- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本指南围绕 Apache Beam Java SDK 的Group变换展开,讲解如何在不显式构造 key 的情况下,直接按输入 schema 的一个或多个字段对PCollection分组,并对每个分组叠加多个CombineFn聚合(求和、Top N、近似分位数等)。读者学完后,将能掌握Group.byFieldNames、Group.globally及aggregateField构建器的完整用法,理解其输出 schema 的生成规则,并能在 Tour of Beam 的 Playground 练习中直接运行验证。
一、Group 变换:Schema 驱动的分组与聚合
在 Beam 的 Schema 体系中(详见 schema-concept/creating-schema),数据以PCollection<Row>或带 schema 的 POJO 形式流动。Group变换(实现位于 sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/transforms/Group.java)正是针对这类数据设计的通用分组变换:
- 它按输入 schema 中的一个或多个字段对
PCollection中的记录进行分组; - 你还可以对分组结果应用聚合,这恰恰是
Group变换最常见的用途; - 输出结果自带 schema,输出 schema 中的每一个字段对应一个聚合。
当不使用 combiner(聚合函数)时,Group变换的作用与GroupByKey完全一致,但有一个关键差别:你无需显式地提取 key——只需声明按哪些字段分组,Beam 会自动完成 key 的提取与组合。
例如,考虑如下输入类型(带@DefaultSchema(JavaFieldSchema.class)注解的 POJO):
@DefaultSchema(JavaFieldSchema.class) public class UserPurchase { public String userId; public String country; public long cost; public double transactionDuration; }四个字段分别表示用户 ID、国家、消费金额与交易时长,后续所有示例都基于该输入。
二、按字段分组:Group by fields
按用户和国家两个字段对全部购买记录分组,只需一行:
PCollection<Row> byUser = input.apply(Group.byFieldNames("userId", "country"));这里Group.byFieldNames(String... fieldNames)返回一个ByFields<InputT>变换(源码见 Group.java),它通过FieldAccessDescriptor.withFieldNames(fieldNames)描述要提取的字段。
输出 schema 结构
从源码中ByFields.expand的实现(Group.java)可以看到,输出 schema 由两个字段组成:
- key 字段(默认名为
key):类型为Row,包含所有被选为分组依据的字段(本例即userId与country); - value 字段(默认名为
value):类型为Iterable<Row>,包含原始完整行。
底层执行链路为:Convert.toRows()把元素统一转成Row→WithKeys.of用 RowSelector 按声明的字段提取并设置 key(Group.java)→GroupByKey.create()完成标准按键分组 →ParDo组装成输出 Row。
这与手写GroupByKey+WithKeys的效果等价,但代码量和出错概率显著降低。
分组字段的多种指定方式
除了按字段名,Group还提供其他入口(Group.java):
| 静态方法 | 说明 |
|---|---|
byFieldNames(String... fieldNames) | 按字段名分组,最常用 |
byFieldNames(Iterable<String> fieldNames) | 按字段名集合分组 |
byFieldIds(Integer... fieldIds) | 按字段序号分组(schema 字段的 0-based 序号) |
byFieldAccessDescriptor(FieldAccessDescriptor) | 按完整的字段访问描述符分组,最灵活 |
同时可通过withKeyField(String)/withValueField(String)(Group.java)重命名输出中的 key/value 字段。
多字段与嵌套字段分组
多字段分组(如同时按field1与field2)在单元测试 GroupTest.java 中验证:输出 key 字段是一个包含这两个字段的嵌套 Row,value 字段按组合 key 聚合对应记录。
FieldAccessDescriptor还支持点号路径,因此可以按嵌套字段分组,例如Group.byFieldNames("inner.field1", "inner.field2"),对应测试 testGroupByNestedKey。
三、分组聚合:Group with aggregation
实际业务中,分组的目的几乎都是为了聚合。Group类内部的构建器方法允许你为输入 schema 的每个字段(或字段集合)分别创建聚合,并据此自动生成输出 schema。例如:
PCollection<Row> aggregated = input .apply(Group.byFieldNames("userId", "country") .aggregateField("cost", Sum.ofLongs(), "total_cost") .aggregateField("cost", Top.<Long>largestLongsFn(10), "top_purchases") .aggregateField("cost", ApproximateQuantilesCombineFn.create(21), Field.of("transactionDurations", FieldType.array(FieldType.INT64)));结果将是一个包含total_cost、top_purchases、transactionDurations三个字段的新 Row schema:
- total_cost:该用户在该国家的所有购买金额之和(
Sum.ofLongs()); - top_purchases:金额最高的前十笔购买(
Top.largestLongsFn(10),输出为数组字段); - transactionDurations:交易时长的直方图(
ApproximateQuantilesCombineFn.create(21)产出 21 个分位数点,输出为INT64数组)。
同时,输出 schema 中还包含一个key 字段,它是一个嵌套 Row,内含userId和country。测试 testByKeyWithSchemaAggregateFn 展示了一个典型结果:key 字段存分组依据,value 字段存所有聚合结果,二者均为 Row 类型。
输出字段类型推断与 Java 类型擦除
通常,字段类型可以从传入的
Combine.CombineFn自动推断出来。然而,由于 Java 的类型擦除(type erasure),有时无法推断。这种情况下,你需要用Schema.Field显式指定字段类型。上面的例子中,transactionDurations字段就是显式指定类型的。
从源码看,aggregateField有两种重载:传入String outputFieldName时(Group.java)由框架从CombineFn的泛型签名推断输出类型;传入Field outputField时(Group.java)则直接使用你给定的字段定义。当使用ApproximateQuantilesCombineFn这类中间类型信息不足的CombineFn时,务必走第二种重载。
底层聚合原理:CombineFieldsByFields
带聚合的分组由CombineFieldsByFields完成(Group.java):
ToKvs:按声明的字段把输入转成KV<Row, Row>;Combine:对每个 key 执行聚合。getCombineTransform(Group.java)会根据场景自动选择Combine.fewKeys或Combine.perKey(默认倾向于fewKeys,因为通常只挑选少数字段,combiner 提升更有利);ToRow:把KV组装回带 schema 的 Row。
聚合函数本身通过SchemaAggregateFn组合多个CombineFn并自动推导输出 schema,所有aggregateField/aggregateFields调用的并集决定最终输出 schema。
分组聚合的性能调优选项
CombineFieldsByFields还暴露了两个构建器方法(Group.java):
withPrecombining(boolean value):预合并开关,默认开启。当唯一 key 很多、且 combiner 中间状态比平均行更大时,预合并反而有害,此时可关闭;withHotKeyFanout(int n)/withHotKeyFanout(SerializableFunction<Row, Integer> f):针对热 key 设置扇出,把单个 key 的聚合拆成 n 份并行,缓解数据倾斜。
对应的全局聚合版本(CombineFieldsGlobally)也提供withFanout(int)(Group.java)。
四、全局分组与全局聚合:Group.globally
当不需要按字段分组、而是要对整个PCollection做聚合时,使用Group.globally()(Group.java):
// 不带 combiner:把所有元素聚成一个 Iterable PCollection<Iterable<Basic>> all = pipeline.apply(Create.of(elements)) .apply(Group.globally()); // 带 combiner:全局计数 Group.CombineGlobally<Basic, Long> count = Group.<Basic>globally().aggregate(Count.combineFn()); PCollection<Long> total = pipeline.apply(Create.of(elements)).apply(count);全局分组在源码中通过WithKeys.of((Void) null)附加 null key →GroupByKey→Values实现(Group.java),本质与Combine.globally一致。全局字段聚合示例见测试 testGloballyWithSchemaAggregateFn。
五、进阶用法:多字段聚合、Base 值聚合与字段名组合
- 多字段聚合:
aggregateFields(List<String> inputFieldNames, CombineFn fn, ...)(Group.java)允许把一个CombineFn同时作用于多个输入字段;对应的aggregateFieldsById则按字段序号指定。 - Base 值聚合:
aggregateFieldBaseValue(...)系列方法(如 Group.java)对字段的基础值而非原始值做聚合,适用于枚举、逻辑类型等场景——测试 GroupTest.java 即对枚举字段求和验证了该语义。 - 替换聚合函数:
aggregateField的第三个参数是输出字段名,因此可自由更换CombineFn。比如把Sum换成Max:
.apply(Group.byFieldNames("userName").aggregateField("score", Max.ofIntegers(), "total"))六、Playground 实战:游戏用户统计
在 Tour of Beam 的 Playground 窗口中(learning/tour-of-beam/learning-content/schema-based-transforms/group/java-example/Task.java),可以直接运行Group的完整示例。该示例从gs://apache-beam-samples/game/small/gaming_data.csv读取游戏数据(字段为userId,userName,score,gameId,date),随机抽样 100 条后按用户分组:
PCollection<String> pCollection = input .apply(MapElements.into(TypeDescriptor.of(Object.class)).via(it -> it)) .setSchema(type, ...) .apply(Group.byFieldNames("userId").aggregateField("score", Sum.ofIntegers(), "total")) .apply(MapElements.into(TypeDescriptor.of(String.class)) .via(row -> row.getRow(0).getValue(0) + " : " + row.getRow(1).getValue(0)));运行后,你将看到每个userId的累计得分(score之和),即"某个游戏中的用户统计"。这里的关键阅读点:
setSchema(...)先把 POJO 流显式声明为含 5 个字段的 schema;Group.byFieldNames("userId")按用户分组,aggregateField("score", Sum.ofIntegers(), "total")对每组得分求和;- 输出 Row 中
getRow(0)是 key 字段(内含userId),getRow(1)是聚合结果字段(内含total),示例据此拼接可读字符串。
示例的元信息(unit-info.yaml)标注其复杂度为ADVANCED,属于 schema-based-transforms 学习路径中的进阶单元。整个练习覆盖了本指南第二、三节的全部核心 API。
七、测试验证:行为与输出 schema 的双重保障
仓库 sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/transforms/GroupTest.java 提供了覆盖上述所有行为的单元测试:
| 测试方法 | 验证内容 |
|---|---|
testGroupByOneField(L103) | 单字段分组,key/value 输出 schema 与预期完全一致 |
testGroupByMultiple(L147) | 多字段分组,组合 key 的 Row 结构 |
testGroupByNestedKey(L205) | 点号路径按嵌套字段分组 |
testGroupGlobally(L249) | 无 combiner 的全局分组输出 Iterable |
testGlobalAggregationWithFanout(L265/L275) | 全局聚合及withFanout行为 |
testByKeyWithSchemaAggregateFn(L331) | 按键聚合多字段(Sum + Top 组合) |
testGloballyWithSchemaAggregateFn(L375) | 全局多字段聚合,输出类型(INT64、INT32、数组)自动推导 |
这些测试均标注@Category(NeedsRunner.class),可在 Direct Runner 上实际运行,是理解Group输出 schema 推导规则(key 字段 + 聚合字段并集)最直接的参考资料。
八、小结
Group变换把 Beam 中"提取 key →GroupByKey→ 自定义聚合 → 组装结果"的冗长链路压缩为一条声明式 API:byFieldNames声明分组依据,aggregateField/aggregateFields声明聚合逻辑,输出 schema 自动生成。它的适用场景包括:
- 按一个或多个字段(含嵌套字段)分组,等价于免 key 的
GroupByKey; - 对每个分组叠加求和、Top N、近似分位数、最大值等多种
CombineFn; - 不关心 key 时用
Group.globally()做全局聚合。
值得记住的两个要点:一是遇到ApproximateQuantilesCombineFn等因 Java 类型擦除而无法推断输出类型的CombineFn时,必须用Schema.Field显式声明字段类型;二是热 key 倾斜时可通过withHotKeyFanout、withPrecombining(false)进行调优。在 Tour of Beam 的 Playground 中运行游戏统计示例,即可完整验证本文所有结论。
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
基于 Schema 的 Group 变换:Apache Beam 按字段分组与聚合的完整指南
基于 Schema 的 Group 变换:Apache Beam 按字段分组与聚合的完整指南 导读 :本文围绕 Apache Beam Java SDK 中面向
大数据批处理流处理数据工程Apache Beam Java 示例实战:用 Beam SQL 与 Schema Transforms 计算按键聚合指标
Apache Beam Java 示例实战:用 Beam SQL 与 Schema Transforms 计算按键聚合指标 本文基于 Apache Beam 仓
大数据批处理流处理数据工程FAST Frame 的 neutralLayerFloating 设计令牌解析:为 Flyout / Menu 等浮动层构建自适应配色
FAST Frame 的 neutralLayerFloating 设计令牌解析:为 Flyout / Menu 等浮动层构建自适应配色 导读 neutralL
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考