ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam 模式化 Join:基于 Schema 的等值连接变换实战指南

Apache Beam 模式化 Join:基于 Schema 的等值连接变换实战指南 批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载本文围绕 Apache Beam 中基于 Schema 的等值连接Equi-join变换Join展开讲解如何对两个带 Schema 的PCollection按字段执行内连接、外连接等六类连接操作并结合仓库源码、单元测试与 Playground 示例帮助读者掌握从自然连接到异构字段名连接的完整实战用法。关联文档join/description.md一、什么是 Schema 化 JoinJoin是 Apache Beam Java SDK 中位于org.apache.beam.sdk.schemas.transforms包下的一个变换专门用于对两个带 Schema 的PCollection执行等值连接equijoin。与传统的CoGroupByKey相比它最大的价值在于用户无需手工抽取键、构造 KV只需声明「按哪些字段连接」其余工作按键分组、字段解析、结果拼装全部由变换自动完成。核心 API 定义在源码 Join.java 中其类注释明确了变换语义This transform allows joins between two input PCollections simply by specifying the fields to join on.连接结果是一个PCollectionRow其中包含两个字段名称分别为lhs和rhs分别承载左、右输入PCollection的完整行数据。这两个标签在源码中以常量形式暴露Join.javapublic static final String LHS_TAG lhs; public static final String RHS_TAG rhs;也就是说连接的产物是一个二元组式 RowRow { lhs: Row(左表整行), rhs: Row(右表整行) }两个子 Row 各自保持输入侧的原 Schema。二、前置条件PCollection 必须带 SchemaJoin只适用于带 Schema 的PCollection即PCollectionRow或通过setSchema/DefaultSchema注解绑定 Schema 的 POJO 类型。在 Playground 示例 Task.java 中可以看到两种典型做法方式一POJO 注解绑定 SchemaDefaultSchema(JavaFieldSchema.class) public static class Game { public String userId; public Integer score; public String gameId; public String date; SchemaCreate public Game(String userId, Integer score, String gameId, String date) { ... } }方式二显式构造 Schema 并 setSchemaSchema gameSchema Schema.builder() .addStringField(userId) .addInt32Field(score) .addStringField(gameId) .addStringField(date) .build(); PCollectionGame gameInfo getGamePCollection(pipeline) .setSchema(gameSchema, TypeDescriptor.of(Game.class), toRowFn, fromRowFn);只有两侧PCollection都具备 SchemaJoin才能解析字段名、自动完成编码与行转换。三、自然连接using(字段名)当左右两个PCollection拥有同名同类型的连接字段时直接使用using(...)指定字段名即可完成「自然连接」Natural Join。这是原文档给出的第一个示例PCollectionRow joined input1.apply(Join.innerJoin(input2).using(user, country));该调用会在两侧 Schema 上同时解析user与country两个字段并作为连接键。从源码实现看using方法实际上就是把同一组字段名同时传给左右两侧的FieldAccessDescriptorJoin.javapublic Join.ImplLhsT, RhsT using(String... fieldNames) { return new Join.Impl(joinType, rhs, FieldsEqual.left(fieldNames).right(fieldNames)); }除字符串字段名外using还提供了两个重载using(Integer... fieldIds)按字段在 Schema 中的索引位置连接using(FieldAccessDescriptor)按字段访问描述符连接。四、异构字段名连接on(FieldsEqual.left(...).right(...))如果右侧PCollection的连接字段与左侧名称不同但类型必须匹配则需要使用on(...)结合FieldsEqual分别指定两侧字段。这是原文档给出的第二个示例PCollectionRow joined input1.apply(Join.innerJoin(input2) .on(FieldsEqual.left(user, country).right(otherUser, otherCountry)));其中FieldsEqual是一个「谓词对象」源码Join.java为其提供了三组对称的重载方法参数类型说明left(String... fieldNames)/right(String... fieldNames)字符串按字段名指定连接字段left(Integer... fieldIds)/right(Integer... fieldIds)整数按 Schema 字段索引指定left(FieldAccessDescriptor)/right(FieldAccessDescriptor)描述符高级用法按字段访问描述符指定FieldsEqual的链式构建支持任意顺序既可以从left(...)开头再.right(...)也可以先right(...)再.left(...)。连接执行前expand方法会调用predicate.resolve(lhsSchema, rhsSchema)把字段名解析为具体的字段引用Join.java因此字段名校验发生在执行期。五、支持的连接类型Supported methods原文档列出的六种连接类型均可通过替换工厂方法名快速切换连接类型工厂方法语义内连接Inner joinJoin.innerJoin(rhs)仅保留两侧都能匹配上的行全外连接Full outer joinJoin.fullOuterJoin(rhs)保留两侧全部行缺失侧置 null左外连接Left outer joinJoin.leftOuterJoin(rhs)保留左侧全部行右侧缺失置 null右外连接Right outer joinJoin.rightOuterJoin(rhs)保留右侧全部行左侧缺失置 null左内连接Left inner join交换左右输入后使用innerJoin以左侧为基准的内连接右内连接Right inner join交换左右输入后使用innerJoin以右侧为基准的内连接从源码看Join.javaJoin类还额外提供了两种广播Broadcast变体innerBroadcastJoin(rhs)内连接右侧以 SideInput 广播方式参与leftOuterBroadcastJoin(rhs)左外连接右侧以 SideInput 广播方式参与。这两者在expand中通过withSideInput()标记实现Join.java适用于右侧数据集较小、希望通过广播避免 shuffle 的场景。外连接的结果可空性语义从单元测试 JoinTest.java 可以确认各类连接在输出 Schema 上的差异内连接lhs与rhs均为非空 Row 字段addRowField左外连接lhs非空、rhs可空addNullableField未匹配行 rhs 为 nullJoinTest.java右外连接lhs可空、rhs非空未匹配行 lhs 为 nullJoinTest.java全外连接lhs、rhs均可空JoinTest.java。这些语义与 SQL 外连接完全一致便于迁移 SQL 经验。六、底层实现基于 CoGroup 与 FieldAccessDescriptorJoin并不是一个从零实现连接逻辑的变换而是构建在 Schema 化CoGroup之上的高层封装。查看expand方法Join.java可以看到其实现策略将左右两个PCollection放入PCollectionTuple分别以lhs、rhs为 tag根据JoinType构造CoGroup.join(...)链内连接两侧都按字段连接并调用crossProductJoin()展开笛卡尔积外连接在对应侧追加withOptionalParticipation()表示该侧可「缺席」缺席即补 null广播变体追加withSideInput()将右侧作为 SideInput 广播最终由CoGroup的crossProductJoin()产出PCollectionRow。字段解析完全依赖 Schema 层的FieldAccessDescriptor——它负责把「字段名 / 字段索引」解析为对 Schema 的实际引用这也是Join能天然支持 POJO、Row 等不同类型输入的原因。正因为Join基于CoGroup两者共享同一套底层机制当需要连接三个及以上的 PCollection 时可直接使用CoGroup见 co-group/description.md例如PCollectionRow input PCollectionTuple.of(input1, input1, input2, input2, input3, input3) .apply(CoGroup.join(CoGroup.By.fieldNames(user, country)));七、完整可运行示例游戏用户与得分统计Playground 目录下的 Task.java 提供了一个端到端可运行示例元信息见 unit-info.yaml复杂度标记为 ADVANCED数据准备从公共数据集gs://apache-beam-samples/game/small/gaming_data.csv读取 CSV采样 100 行分别解析为UseruserId、userName与GameuserId、score、gameId、date两类带 Schema 的 PCollection。执行连接PCollectionRow pCollection userInfo.apply( Join.User, GameinnerJoin(gameInfo).using(userId));输出查看对连接结果应用ParDo打印每一行pCollection.apply(User flatten row, ParDo.of(new LogOutput(Flattened)));运行后日志中的每一行将形如Flattened: Row:[lhs: Row:[userId: ..., userName: ...], rhs: Row:[userId: ..., score: ..., gameId: ..., date: ...]]可以看到userId同时出现在lhs与rhs中——这正是自然连接的特征键字段本身不会在顶层被抽离而是完整保留在两侧行内。八、Playground 练习一行代码切换连接类型原文档的 Playground 练习明确指出仅需修改方法名即可切换连接类型。将示例中的内连接换成全外连接.apply(Join.fullOuterJoin(gameInfo).using(userId));运行后所有在gameInfo中没有对应userId的User行也会出现在结果中此时rhs字段为 null。类似地可以逐一尝试.apply(Join.leftOuterJoin(gameInfo).using(userId)); // 左外连接 .apply(Join.rightOuterJoin(gameInfo).using(userId)); // 右外连接 .apply(Join.innerBroadcastJoin(gameInfo).using(userId)); // 广播内连接九、测试验证与正确性保障Join的语义正确性由 JoinTest.java 中七组Category(NeedsRunner.class)测试覆盖全部基于TestPipelinePAssert断言包括testInnerJoinSameKeys/testInnerJoinDifferentKeys验证自然连接与异构字段名连接产出相同结果并校验输出 SchematestInnerJoinDifferentKeysNullable2验证右表存在 null 键时的内连接行为testOuterJoinSameKeys/testOuterJoinDifferentKeys验证全外连接补 null 行为testLeftOuterJoinSameKeys/testRightOuterJoinSameKeys验证单侧外连接的可空字段语义。这些测试同时断言了Join的输出 Schema 结构例如内连接的期望 Schema 为Schema.builder() .addRowField(Join.LHS_TAG, CG_SCHEMA_1) .addRowField(Join.RHS_TAG, CG_SCHEMA_1) .build();如果你在自己的代码中手动构造期望 Schema可直接复用Join.LHS_TAG/Join.RHS_TAG常量避免硬编码字符串。十、小结与进一步学习Join把「两表等值连接」这件事压缩为一行链式调用语法上非常接近 SQL 直觉Join.innerJoin(rhs).using(key) // 自然连接 Join.leftOuterJoin(rhs).on(FieldsEqual.left(a).right(b)) // 异构键左外连接其背后依赖的是 Beam 的 Schema 体系Schema、Row、FieldAccessDescriptor与 Schema 化CoGroup的联合实现在 Java SDK 的schemas.transforms包内与Select、Group、Filter、Convert、Rename等变换共同构成完整的 Schema 变换家族见 module-info.yaml。若需多表连接、自定义键输出结构建议继续研读 CoGroup 文档若需深入源码可重点阅读 Join.java 的expand与FieldsEqual实现以及 CoGroup.java 的crossProductJoin逻辑。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐ComfyUI-WanVideoWrapper小显存跑 WanVideo 新模型的完整实战指南ComfyUI WanVideoWrapper小显存跑 WanVideo 新模型的完整实战指南 显存不够、Wan 系模型更新又快刚追上一个版本下一代已经放Apache Beam 基于 Schema 的多 PCollection 等值连接CoGroup 变换完整实战指南Apache Beam 基于 Schema 的多 PCollection 等值连接CoGroup 变换完整实战指南 Apache Beam 的 CoGroup大数据批处理流处理数据工程Apache Beam Schema-Based Joins 实战用 Join transform 在带 Schema 的 PCollection 上执行等值连接Apache Beam Schema Based Joins 实战用 Join transform 在带 Schema 的 PCollection 上执行等值大数据批处理流处理数据工程上一篇使用 Agentic Actions Auditor 审计 GitHub Actions 中的 AI Agent 集成安全下一篇NetBox 配置上下文Config Context完全指南作用域分配、分层合并与预渲染缓存原理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进