news 2026/10/10 5:59:48

Apache Beam Schema 化 Group 变换:按字段分组与多聚合实战指南(Java SDK)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Schema 化 Group 变换:按字段分组与多聚合实战指南(Java SDK)
  • 批处理
  • 流处理
  • 大数据

【免费下载链接】beam

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

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

本指南围绕 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):

  1. ToKvs:按声明的字段把输入转成KV<Row, Row>;
  2. Combine:对每个 key 执行聚合。getCombineTransform(Group.java)会根据场景自动选择Combine.fewKeys或Combine.perKey(默认倾向于fewKeys,因为通常只挑选少数字段,combiner 提升更有利);
  3. 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.

项目地址:https://gitcode.com/gh_mirrors/beam15/beam
点击查看免费下载
上一篇:LeetCode 109 有序链表转换二叉搜索树:双指针中点递归与数组缓存两种构建方案详解
下一篇:awesome-math 完整数学资源清单:从入门自学到形式化证明的选型指南

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/10 5:56:46

flet-desktop 解析:Flutter 桌面客户端如何把 Flet 应用变成原生窗口

前端跨平台桌面应用移动开发 【免费下载链接】flet Build realtime web, mobile and desktop apps in Python only. No frontend experience required. 项目地址&#xff1a; https://gitcode.com/gh_mirrors/fl/flet 点击查看 免费下载 导读 flet-desktop 是 Flet 生态中专用…

作者头像 李华
网站建设 2026/10/10 5:54:18

strace高级技巧与生产实战:从系统调用定位线上问题

上周线上有个服务接口偶发超时&#xff0c;日志里只能看到“上游组件超时”&#xff0c;CPU、内存全正常&#xff0c;看监控也找不到异常。我挂上 strace 抓了不到 10 分钟&#xff0c;就从系统调用时间戳里找到了真正的等待点。群里一个同事问了一句&#xff1a;“strace 还能…

作者头像 李华
网站建设 2026/10/10 5:53:00

kshell:为散落一地的AI编程会话建一个本地中央车站

装了一堆 AI 编程工具之后&#xff0c;真正让人抓狂的已经不是“哪个更好用”&#xff0c;而是会话散落一地——今天这个问题是在工具 A 里问的&#xff0c;那个报错是在工具 B 里解决的&#xff0c;一周以后想翻记录&#xff0c;手忙脚乱也找不到。我自己被这个状态折磨了快两…

作者头像 李华
网站建设 2026/10/10 5:51:48

物联网平台源码实战:从MQTT协议选型到海康摄像头接入

做物联网平台源码这类项目&#xff0c;最容易被低估的其实不是业务功能&#xff0c;而是设备接入层的通信协议——TCP/IP、MQTT、HTTP三条链路怎么分工&#xff0c;海康摄像头怎么取流&#xff0c;传感器报文怎么从一堆字节里把有效数据抠出来&#xff0c;这些东西搞不清楚&…

作者头像 李华
网站建设 2026/10/10 5:50:07

TaoToken 实战:让 AI 帮写注释并直接生成代码的配置指南

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

作者头像 李华
网站建设 2026/10/10 5:49:23

PCA9422+PIC32MX构建可编程电源管理子系统

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

作者头像 李华