资讯动态

Apache Beam Java 实战:用 Side Input 为 ParDo 注入运行时查找数据(Kata:城市→国家映射)

发布时间:2026/10/9 1:40:31 来源:尧图企业网站定制
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载本文以 Apache Beam 官方 Java Katas 中的「Side Input」练习为核心讲解如何在 ParDo 之外构建一个可被 DoFn 按需读取的辅助数据集PCollectionView实现用城市查找国家这类典型的运行时数据注入场景。读完本文你将掌握View.asMap()创建旁路视图、ParDo.withSideInputs()绑定视图、context.sideInput()读取数据的完整调用链并能结合仓库内的源码与单元测试独立完成并验证同类任务。什么是 Side Input在主输入之外注入运行时数据在 Apache Beam 中一个ParDo变换的典型形态是单个主输入 PCollection → 逐元素处理 → 输出但在很多真实业务里处理每个元素时往往还需要参考一份额外的、在运行时才能确定的数据。例如根据订单中的城市 ID 查询城市所属的国家/地区根据用户 ID 查询其会员等级决定是否发放优惠依据当日配置的阈值决定元素是否过滤。Beam 为此提供的机制就是Side Input旁路输入。正如本 Kata 的 task.md 所述除了主输入 PCollection 之外你可以以旁路输入的形式为 ParDo 变换提供额外的输入。旁路输入是 DoFn 在处理主输入 PCollection 的每一个元素时都可以访问到的附加输入。当你指定一个旁路输入时你实际上是创建了另一份数据的视图view这份视图可以在 ParDo 变换的 DoFn 处理每个元素的过程中被读取。Side Input 的核心价值在于这份附加数据不是硬编码在代码里的而是由输入数据本身或管线的另一个分支在运行时动态计算出来的。这也正是它区别于普通全局常量、静态 Map 的关键所在——你可以在同一份管线中让另一个分支先聚合、计算出一份结果再把它作为旁路视图供主分支逐元素查阅。Kata 任务为每位 Person 补全所在国家本练习位于learning/katas/java/Core Transforms/Side Input/Side Input/要求完成如下任务Kata请根据每个人Person所在的城市city为其补充所在的国家country。任务给定的数据如下城市→国家映射旁路数据Beijing→China、London→United Kingdom、San Francisco→United States、Singapore→Singapore、Sydney→Australia人员主输入Henry(Singapore)、Jane(San Francisco)、Lee(Beijing)、John(Sydney)、Alfred(London)期望输出每个Person对象携带name、city、country三个字段其中country由城市查表得出。任务中给出的三个关键提示依次是使用View创建citiesToCountries的PCollectionView使用接受旁路输入的ParDoDoFn即withSideInputs并建议参阅 Beam 官方编程指南中 Side inputs 一节。数据模型Person练习的数据模型定义在 Person.java。它是一个实现了Serializable的 POJO同时提供了两个构造函数——Person(name, city)用于构造未补全国家的输入数据Person(name, city, country)用于构造已补全国家的输出数据并实现了equals/hashCode/toString这使得后续可以用PAssert.containsInAnyOrder直接比较元素内容public class Person implements Serializable { private String name; private String city; private String country; public Person(String name, String city) { ... } public Person(String name, String city, String country) { ... } // getName / getCity / getCountry / equals / hashCode / toString }第一步用 View 把旁路数据变成 PCollectionViewSide Input 不是直接把一个PCollection传给 DoFn而是要先把旁路 PCollection转换成一个视图对象PCollectionView。在 Task.java 中这一步被封装在createView方法里static PCollectionViewMapString, String createView( PCollectionKVString, String citiesToCountries) { return citiesToCountries.apply(View.asMap()); }View.asMap()会把PCollectionKVK, V物化成一个以 K 为键、V 为值的Map视图返回类型为PCollectionViewMapString, String。PCollectionView是org.apache.beam.sdk.values包下的核心类型见 View.java 的工厂方法族它不参与逐元素的数据流而是作为一份可随机访问的物化数据被旁路绑定到具体变换上。Beam 的View工厂还提供了其他常用的视图形态适用于不同的旁路数据场景视图工厂方法适用场景View.asMap()PCollectionKVK, V→MapK, V按键随机查找本 Kata 所用View.asSingleton()单元素 PCollection → 单个值用于全局配置/阈值等View.asList()PCollection →ListT用于全量列表旁路View.asIterable()PCollection →IterableT用于可迭代的全量数据第二步用 ParDo withSideInputs 绑定旁路视图视图创建好之后需要把它挂载到目标ParDo变换上这一步由ParDo.of(...).withSideInputs(view)完成。withSideInputs接受一个或多个PCollectionView因此一个 DoFn 可以同时绑定多个旁路视图static PCollectionPerson applyTransform( PCollectionPerson persons, PCollectionViewMapString, String citiesToCountriesView) { return persons.apply(ParDo.of(new DoFnPerson, Person() { ProcessElement public void processElement(Element Person person, OutputReceiverPerson out, ProcessContext context) { MapString, String citiesToCountries context.sideInput(citiesToCountriesView); String city person.getCity(); String country citiesToCountries.get(city); out.output(new Person(person.getName(), city, country)); } }).withSideInputs(citiesToCountriesView)); }这里需要特别留意 DoFn 方法签名为了读取旁路输入ProcessElement方法必须声明一个ProcessContext context参数或直接声明SideInput PCollectionView参数然后在方法体内通过context.sideInput(citiesToCountriesView)获取已物化的Map即可像使用普通 Java Map 一样按键查询。这是 Side Input 与普通 ParDo 在写法上最核心的差异。第三步组装管线并查看运行输出main方法把整条链路串起来先创建两个分支城市→国家映射、人员列表再经过createView→applyTransform最后用 Kata 工具类Log.ofElements()定义于 Log.java把结果逐条打印到日志随后pipeline.run()启动执行PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline Pipeline.create(options); PCollectionKVString, String citiesToCountries pipeline.apply(Cities and Countries, Create.of( KV.of(Beijing, China), KV.of(London, United Kingdom), KV.of(San Francisco, United States), KV.of(Singapore, Singapore), KV.of(Sydney, Australia) )); PCollectionViewMapString, String citiesToCountriesView createView(citiesToCountries); PCollectionPerson persons pipeline.apply(Persons, Create.of( new Person(Henry, Singapore), new Person(Jane, San Francisco), new Person(Lee, Beijing), new Person(John, Sydney), new Person(Alfred, London) )); PCollectionPerson output applyTransform(persons, citiesToCountriesView); output.apply(Log.ofElements()); pipeline.run();预期运行日志每条 Person 均带有补全后的国家字段Person{nameHenry, citySingapore, countrySingapore} Person{nameJane, citySan Francisco, countryUnited States} Person{nameLee, cityBeijing, countryChina} Person{nameJohn, citySydney, countryAustralia} Person{nameAlfred, cityLondon, countryUnited Kingdom}Log.ofElements()是一个基于ParDo的辅助PTransform见 Log.java它通过 SLF4J 输出每个元素且会原样透传元素不影响后续处理在非全局窗口如固定窗口/滑动窗口场景下它还会在日志中附带窗口信息便于调试窗口化管线。用单元测试验证 Side Input 的正确性Kata 自带完整的 JUnit 测试 TaskTest.java它直接复用Task.createView与Task.applyTransform两个静态方法借助TestPipeline与PAssert对结果做无序全量断言Rule public final transient TestPipeline testPipeline TestPipeline.create(); Test public void sideInput() { // ... 构造 citiesToCountries 与 persons与 main 相同的数据 PCollectionViewMapString, String citiesToCountriesView Task.createView(citiesToCountries); PCollectionPerson results Task.applyTransform(persons, citiesToCountriesView); PAssert.that(results) .containsInAnyOrder( new Person(Henry, Singapore, Singapore), new Person(Jane, San Francisco, United States), new Person(Lee, Beijing, China), new Person(John, Sydney, Australia), new Person(Alfred, London, United Kingdom) ); testPipeline.run().waitUntilFinish(); }该测试验证了两点其一createViewapplyTransform的组合能把每个城市正确映射到国家其二Person的equals实现让PAssert.containsInAnyOrder可以直接按值比较不依赖元素顺序。这是 Beam 单元测试的标准姿势——TestPipeline.create()注册管线生命周期PAssert.that(...)声明期望run().waitUntilFinish()执行并校验。原理小结Side Input 的完整调用链与适用边界从本 Kata 可以提炼出 Side Input 的完整调用链旁路 PCollectionKVK,V └─ View.asMap() → PCollectionViewMapK,V 主输入 PCollectionPerson └─ ParDo.of(DoFn).withSideInputs(view) └─ ProcessElement(..., ProcessContext context) └─ context.sideInput(view) → MapK,V按键查询结合 View.java 与PCollectionView的定义可以看出旁路视图本质上是将一份相对较小的参考数据物化后广播给所有处理主输入元素的并行任务让每个元素在处理时都能本地随机访问。使用时需要注意以下几点适用场景旁路数据应在运行时确定由输入数据或其他分支计算得出且整体规模适合被完整物化与分发若参考数据量极大应评估改为 CoGroupByKey 等流式关联方式是否更合适读取方式必须在 DoFn 的ProcessElement中通过context.sideInput(view)或SideInput参数注入读取不能在构造 DoFn 实例时读取多视图withSideInputs支持一次绑定多个视图一个 DoFn 可同时查阅多份旁路数据视图形态根据旁路数据的形态选择asMap/asSingleton/asList/asIterable本 Kata 的按键查值场景对应asMap。完成本练习后你可以继续在learning/katas/java/Core Transforms/目录下探索相邻主题如 DoFn Additional Parameters、CoGroupByKey、Side Output它们与 Side Input 共同构成了 Beam 核心变换中多路数据协作的完整能力集。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Python Kata 实战用 Side Input 为 ParDo 注入运行时附加数据Apache Beam Python Kata 实战用 Side Input 为 ParDo 注入运行时附加数据 导读 Side Input侧输入是 Ap大数据批处理流处理数据工程Apache Beam Kotlin SDK 实战用 Side InputPCollectionView为 ParDo 注入运行时查找数据Apache Beam Kotlin SDK 实战用 Side InputPCollectionView为 ParDo 注入运行时查找数据 Apache大数据批处理流处理数据工程Apache Beam Go SDK 实战用 Side Input 在 ParDo 中注入运行时辅助数据Side Input Kata 演练Apache Beam Go SDK 实战用 Side Input 在 ParDo 中注入运行时辅助数据Side Input Kata 演练 本文以 Ap上一篇终极指南使用Dio发送GraphQL请求的完整教程下一篇终极指南PHPUnit-Mock-Objects异常处理与边界情况测试策略创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑