news 2026/10/10 9:01:24

Apache Beam Java FlattenWith 变换:用链式合并多个 PCollection 的完整指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Java FlattenWith 变换:用链式合并多个 PCollection 的完整指南
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

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

FlattenWith是 Apache Beam Java SDK 中Flatten系列变换的链式变体,用于将多个同类型PCollection合并(merge)为一个逻辑PCollection。与经典Flatten.pCollections()需要先构造PCollectionList不同,FlattenWith可以直接内联在 transform 链的任意位置,还能直接接收Create、Read等根PCollection产生型变换,本文将从官方文档出发,结合 Flatten.java 源码与 FlattenTest.java 测试,完整讲解其语义、用法、底层原理与注意事项。

什么是 FlattenWith

按照 flattenWith.md 的定义:

Merges multiplePCollectionobjects into a single logicalPCollection. It allows for the combination of both rootPCollection-producing transforms (likeCreateandRead) and existing PCollections.

即:FlattenWith 将多个PCollection合并为一个逻辑PCollection,并且既可以合并已经存在的PCollection,也可以合并Create、Read这类"根PCollection产生型变换"。这是它与普通Flatten相比最突出的能力差异——普通Flatten的输入必须是一个PCollectionList,而FlattenWith的第二个输入可以直接是一个尚未运行的变换(PTransform)。

在 Beam Java SDK 中,FlattenWith并不是一个独立的类,而是Flatten类提供的两个静态重载方法:

  • Flatten.with(PCollection<T> other):把输入PCollection与一个已存在的PCollection合并;
  • Flatten.with(PTransform<PBegin, PCollection<T>> other):把输入PCollection与一个**根变换(如Create.of(...))**合并,该变换会在展开时自动应用到PBegin上并执行。

两种重载的返回类型都是PTransform<PCollection<T>, PCollection<T>>,即接收一个PCollection输出一个PCollection,因此可以无缝地放在.apply(...)链中。

与 Flatten.pCollections() 的关系

从源码可以清楚看到,FlattenWith在语义上完全等价于"先构造PCollectionList再调用pCollections()",只是为了内联链式使用而提供的便捷封装。以第一个重载为例,Flatten.java 中FlattenWithPCollection.expand()的实现为:

public PCollection<T> expand(PCollection<T> input) { return PCollectionList.of(input).and(other).apply(pCollections()); }

即:PCollectionList.of(input).and(other)把输入与other打包成一个包含两个PCollection的列表,然后交给pCollections()完成真正的合并逻辑。其getKindString()返回"Flatten.With",用于在调试与指标中标识该变换。

第二个重载(接收根变换)的expand()则等价于:

public PCollection<T> expand(PCollection<T> input) { return PCollectionList.of(input) .and(input.getPipeline().apply(other)) // 先运行 other 变换得到 PCollection .apply(pCollections()); }

即先在输入所在的 Pipeline 上执行other变换,得到其输出PCollection,再与输入一起走pCollections()合并。

何时使用 FlattenWith

FlattenWith最典型的应用场景是"在数据流链的中间位置动态汇入其他数据源"。Beam 官方编程指南 programming-guide.md 中给出了一个非常直观的链式示例:

PCollection<String> merged = pc1 .apply(...) // 在这一步把 pc2 的元素汇入... .apply(FlattenWith.of(pc2)) .apply(...) // 在这一步把 pc3 的元素汇入... .apply(FlattenWith.of(pc3)) .apply(...);

可以看到,通过链式调用FlattenWith.of(...),可以在管道(pipeline)处理流程的不同阶段分别汇入不同的数据集合,而无需事先把所有集合聚合成一个PCollectionList。这种写法比PCollectionList.of(pc1).and(pc2).and(pc3).apply(Flatten.pCollections())更贴近"流式拼接"的心智模型,也更适合代码重构时按需插入新的数据源。

与之相对,Flatten.pCollections()更适合一次性合并一个已知的PCollectionList的场景。两种方式在结果语义上完全一致,可以根据代码可读性选择。

快速上手:完整可运行的示例

以下是基于 FlattenExample.java 改写、同时演示两种with重载的完整示例。该示例采用了与 Beam Playground 中PG_BEAMDOC_SDK_JAVA_FlattenWith示例一致的Create用法:

import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.Flatten; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.values.PCollection; public class FlattenWithExample { public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.create(); Pipeline pipeline = Pipeline.create(options); // 方式一:Flatten.with(PCollection) —— 合并已存在的 PCollection PCollection<String> pc1 = pipeline.apply(Create.of("Hello")); PCollection<String> pc2 = pipeline.apply(Create.of("World", "Beam")); PCollection<String> pc3 = pipeline.apply(Create.of("Is", "Fun")); PCollection<String> merged = pc1 .apply(Flatten.with(pc2)) // 在 pc1 之后汇入 pc2 .apply(Flatten.with(pc3)); // 再汇入 pc3 // 方式二:Flatten.with(PTransform) —— 直接合并根变换的输出(无需先 apply) PCollection<String> merged2 = pipeline.apply(Create.of("Hello")) .apply(Flatten.with(Create.of("World", "Beam", "Is", "Fun"))); // 打印合并结果 merged.apply(ParDo.of(new LogOutput<>("merged: "))); merged2.apply(ParDo.of(new LogOutput<>("merged2: "))); pipeline.run(); } static class LogOutput<T> extends DoFn<T, T> { private static final Logger LOG = LoggerFactory.getLogger(LogOutput.class); private final String prefix; public LogOutput(String prefix) { this.prefix = prefix; } @ProcessElement public void processElement(ProcessContext c) { LOG.info(prefix + c.element()); c.output(c.element()); } } }

运行后merged中会包含"Hello"、"World"、"Beam"、"Is"、"Fun"共 5 个元素(对应 flatten.md 中描述的合并结果),merged2同样包含这 5 个元素——两个重载在结果上是等价的。该示例对应的 Flatten 基础版本还可在 FlattenExample.java 中查看,它演示了标准的PCollectionList+Flatten.pCollections()写法:

PCollection<String> pc1 = pipeline.apply(Create.of("Hello")); PCollection<String> pc2 = pipeline.apply(Create.of("World", "Beam")); PCollection<String> pc3 = pipeline.apply(Create.of("Is", "Fun")); PCollectionList<String> collections = PCollectionList.of(pc1).and(pc2).and(pc3); PCollection<String> merged = collections.apply(Flatten.pCollections());

源码剖析:FlattenWith 的两种重载与底层展开

下面深入 Flatten.java 源码,剖析两种重载的完整实现。

重载一:Flatten.with(PCollection other)

public static <T> PTransform<PCollection<T>, PCollection<T>> with(PCollection<T> other) { return new FlattenWithPCollection<>(other); } private static class FlattenWithPCollection<T> extends PTransform<PCollection<T>, PCollection<T>> { // 仅在管道构建期需要访问,因此标记 transient,不参与序列化 private final transient PCollection<T> other; public FlattenWithPCollection(PCollection<T> other) { this.other = other; } @Override public PCollection<T> expand(PCollection<T> input) { return PCollectionList.of(input).and(other).apply(pCollections()); } @Override public String getKindString() { return "Flatten.With"; } }

要点:

  • other字段声明为transient,因为PCollection仅在管道构建期有意义,运行时(如 Dataflow 作业序列化)不需要携带;
  • expand()将input与other打包成PCollectionList后委托给pCollections(),因此窗口、触发器等校验逻辑全部复用了Flatten本身的实现;
  • getKindString()返回"Flatten.With",便于在 Runner 日志与图(graph)中区分该变换。

重载二:Flatten.with(PTransform<PBegin, PCollection > other)

public static <T> PTransform<PCollection<T>, PCollection<T>> with( PTransform<PBegin, PCollection<T>> other) { return new PTransform<PCollection<T>, PCollection<T>>() { @Override public PCollection<T> expand(PCollection<T> input) { return PCollectionList.of(input) .and(input.getPipeline().apply(other)) // 在管道上执行根变换 .apply(pCollections()); } @Override public String getKindString() { return "Flatten.With"; } }; }

要点:

  • 这里的other是一个PTransform<PBegin, PCollection<T>>,即典型的根变换:Create.of(...)、TextIO.read()...、KafkaIO.read()...等都属于这一类(它们的输入都是PBegin);
  • 在expand()中,Beam 会先在input.getPipeline()上执行other,得到输出PCollection,再与input合并。这就是文档所说"FlattenWith 可以直接接收根PCollection产生型变换"的底层实现;
  • 正因为如此,你可以写出pipeline.apply(Create.of(...)).apply(Flatten.with(TextIO.read().from(path)))这样的链式代码,无需提前单独apply并保存中间变量。

核心合并逻辑:PCollections.expand()

两种重载最终都委托给Flatten.PCollections.expand()(Flatten.java),它负责:

  1. 窗口与触发器校验:对输入列表中每个PCollection的WindowFn与Trigger两两调用isCompatible()检查兼容性,不兼容时抛出IllegalStateException(见下文"窗口兼容性"一节);
  2. 有界性传播:通过IsBounded.BOUNDED.and(input.isBounded())聚合所有输入的界属性——只要任一输入是无界(unbounded)流,输出即为无界;
  3. Coder 选择:默认采用第一个PCollection的 Coder 作为输出 Coder(若列表为空则暂不指定)。

这解释了为什么 flatten.md 中说明"输出 PCollection 的默认 coder 与输入PCollectionList中第一个 PCollection 的 coder 相同,但各输入可以使用不同的 coder,只要元素类型一致"。

测试用例验证:FlattenWith 的实际行为

FlattenTest.java 中针对FlattenWith专门有两个NeedsRunner级别的测试,分别覆盖两种重载:

测试一:testFlattenWithPCollection(FlattenTest.java)

@Test @Category(NeedsRunner.class) public void testFlattenWithPCollection() { PCollection<String> output = p.apply(Create.of(LINES)) .apply("FlattenWithLines1", Flatten.with(p.apply("Create1", Create.of(LINES)))) .apply("FlattenWithLines2", Flatten.with(p.apply("Create2", Create.of(LINES2)))); PAssert.that(output).containsInAnyOrder(flattenLists(Arrays.asList(LINES, LINES2, LINES))); p.run(); }

验证要点:链式调用两次Flatten.with(PCollection),且每次汇入的PCollection本身就是管道上另一个变换的输出(p.apply("Create1", ...)),最终输出包含LINES ∪ LINES2 ∪ LINES的全部元素(containsInAnyOrder说明合并不保证元素顺序)。

测试二:testFlattenWithPTransform(FlattenTest.java)

@Test @Category(NeedsRunner.class) public void testFlattenWithPTransform() { PCollection<String> output = p.apply(Create.of(LINES)) .apply("Create1", Flatten.with(Create.of(LINES))) .apply("Create2", Flatten.with(Create.of(LINES2))); PAssert.that(output).containsInAnyOrder(flattenLists(Arrays.asList(LINES, LINES2, LINES))); p.run(); }

验证要点:直接向Flatten.with(...)传入Create.of(...)这类根变换,Beam 会在展开时自动执行它,输出同样包含三份输入的并集。这正是文档所述"FlattenWith 可以组合根PCollection产生型变换与已存在的 PCollection"的官方测试佐证。

两个测试都用PAssert.that(output).containsInAnyOrder(...)断言,从侧面印证:FlattenWith 合并后的PCollection元素无序,测试不应依赖元素相对顺序。

窗口(Windowing)兼容性:最需要警惕的约束

FlattenWith与普通Flatten共享同一套窗口约束。官方 flatten.md 明确警告:

当使用Flatten合并已应用窗口策略的PCollection时,所有待合并的PCollection必须使用兼容的窗口策略和窗口尺寸。例如,所有被合并的集合都必须(假设性地)使用完全一致的 5 分钟固定窗口,或每 30 秒启动一次的 4 分钟滑动窗口。如果管道试图用Flatten合并窗口不兼容的PCollection,Beam 会在管道构建时抛出IllegalStateException。

从 Flatten.java 的PCollections.expand()可以看到该校验的具体实现:

WindowingStrategy<?, ?> windowingStrategy = inputs.get(0).getWindowingStrategy(); for (PCollection<?> input : inputs.getAll()) { WindowingStrategy<?, ?> other = input.getWindowingStrategy(); if (!windowingStrategy.getWindowFn().isCompatible(other.getWindowFn())) { throw new IllegalStateException( "Inputs to Flatten had incompatible window windowFns: " + windowingStrategy.getWindowFn() + ", " + other.getWindowFn()); } if (!windowingStrategy.getTrigger().isCompatible(other.getTrigger())) { throw new IllegalStateException( "Inputs to Flatten had incompatible triggers: " + windowingStrategy.getTrigger() + ", " + other.getTrigger()); } }

也就是说,合并前 Beam 会逐对比较每个输入的:

  • WindowFn(窗口函数),如FixedWindows、SlidingWindows、Sessions等,通过isCompatible()判断是否兼容;
  • Trigger(触发器),同样要求兼容。

任一不兼容都会在pipeline.run()之前的图构建阶段抛出IllegalStateException,并明确指出是哪两个 WindowFn 或 Trigger 冲突。因此,如果多个数据源采用了不同的窗口策略,需要先通过Window.into(...)统一窗口,再执行FlattenWith。

输出PCollection的窗口策略取自第一个输入集合,且输出元素保留各自输入元素原有的窗口归属与时间戳(见 Flatten.java 的 Javadoc 说明)。

与其他相关变换的对比

flattenWith.md 的 "Related transforms" 一节给出了两个直接相关的变换:

  • Flatten:将多个PCollection合并为一个逻辑PCollection,适合处理多个同类型数据集合的场景。它需要先把输入组织成PCollectionList(PCollectionList.of(pc1).and(pc2)...),是FlattenWith的"批量版";
  • FlatMap:对集合中的每个元素应用一个简单的 1 对多映射函数,每个输入元素可能产生零个或多个输出。它改变的是元素数量与形态,而FlattenWith只改变集合的组织方式(多个集合并成一个集合),不触碰元素本身。

可以这样记忆:Flatten家族做的是"集合层面的合并"(N 个集合 → 1 个集合),FlatMap做的是"元素层面的展开"(1 个元素 → 0..N 个元素)。在实际管道中两者常常配合使用——先用FlattenWith汇拢多个数据源,再统一交给ParDo/FlatMap做后续处理。

小结:FlattenWith 使用速查

场景推荐写法
合并若干已存在的PCollection,一次性完成PCollectionList.of(pc1).and(pc2).apply(Flatten.pCollections())
在 transform 链中间动态汇入某个已存在的PCollectioninput.apply(Flatten.with(pc2))
在 transform 链中间直接汇入Create/Read等根变换input.apply(Flatten.with(Create.of(...)))或input.apply(Flatten.with(TextIO.read().from(path)))
需要合并窗口/触发器兼容的多个流式或批量数据源先统一窗口,再使用上述任一方式合并

使用FlattenWith时请记住三条核心约束:

  1. 元素类型必须一致:所有参与合并的PCollection必须存储相同的数据类型(可使用不同 Coder,但元素类型要一致);
  2. 窗口策略必须兼容:WindowFn 与 Trigger 不兼容会在构建期触发IllegalStateException;
  3. 合并结果无序:输出PCollection只是逻辑上的集合,元素顺序不受保证,测试请使用containsInAnyOrder断言。

如需进一步验证行为,可以直接运行 FlattenTest.java 中的testFlattenWithPCollection与testFlattenWithPTransform两个用例,它们是FlattenWith两种重载最权威的行为参照。

  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载
上一篇:OI-wiki 排序的用法指南:数据预处理、复杂度优化与二分查找实战
下一篇:MLflow × CrewAI 自动追踪(Auto Tracing)实战指南:让多智能体工作流全程可观测

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

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

Midjourney V6风格参考参数--sref详解:原理、调参与实战

用Midjourney V6出图的人&#xff0c;应该都有过这种经历&#xff1a;好不容易磨出一张满意的图&#xff0c;换个构图重新生成&#xff0c;风格却完全跑偏&#xff0c;同样的关键词在不同批次里出来的效果像换了个画师。直到V6版本把风格控制单独拎出来做成一个独立参数&#x…

作者头像 李华
网站建设 2026/10/10 9:00:29

VeapAI 实战(十二)多流程联动:公共服务清单与驱动规则

摘要&#xff1a;本文介绍 VeapAI 如何把「一个流程驱动另一个流程」的联动逻辑从散落的埋点代码&#xff0c;收敛为可视化配置。核心是两张表&#xff1a;公共服务清单&#xff08;wf_public_service&#xff09;声明可复用的联动能力&#xff0c;驱动规则&#xff08;wf_proc…

作者头像 李华
网站建设 2026/10/10 8:59:02

C语言实现小型编译器:从词法分析到四元式的完整指南

简介&#xff1a;面向编译原理课程设计与C语言进阶实践&#xff0c;这份资源以小型编译程序的完整实现为主线&#xff0c;覆盖词法分析、语法分析、语义分析与四元式生成等核心环节&#xff0c;适合计算机专业学生或需要动手理解编译过程的开发者参考。压缩包共4个文件&#xf…

作者头像 李华
网站建设 2026/10/10 8:58:32

Spring Bean实例化的四种方式:从构造器到FactoryBean

很多人学Spring都是从IoC容器开始的&#xff0c;但等你打开一个真实的老项目&#xff0c;看到XML里散落着一堆<bean>配置&#xff0c;有时还是会犯迷糊&#xff1a;同样是创建一个对象&#xff0c;为什么有的直接写个class就完事&#xff0c;有的要加factory-method&…

作者头像 李华
网站建设 2026/10/10 8:55:21

办公用品直售推荐系统全栈实战:SpringBoot+Vue前后端分离设计

1. 项目概述1.1 核心需求解析办公用品直售系统这个方向其实一直很有搞头。大多数企业采购日常办公用品的流程还很原始——行政翻商品目录、人工比价、邮件审批、月底对账&#xff0c;效率低不说&#xff0c;采购记录还不透明。平时咱们做管理系统做得多了&#xff0c;这次我决定…

作者头像 李华
网站建设 2026/10/10 8:55:18

Slack自主AI代理实战:从被动问答到主动处理团队工作流

不少人应该有过这种体验&#xff1a;团队里的 Slack 群聊永远是红点轰炸现场&#xff0c;问个问题没人理、催个进度半天没回音、跨部门协作更是像在玩拼图。大多数团队部署的 AI 机器人&#xff0c;本质是个“问答盒子”&#xff0c;你问一句它答一句&#xff0c;你不问它绝不开…

作者头像 李华