- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
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 multiple
PCollectionobjects 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),它负责:
- 窗口与触发器校验:对输入列表中每个
PCollection的WindowFn与Trigger两两调用isCompatible()检查兼容性,不兼容时抛出IllegalStateException(见下文"窗口兼容性"一节); - 有界性传播:通过
IsBounded.BOUNDED.and(input.isBounded())聚合所有输入的界属性——只要任一输入是无界(unbounded)流,输出即为无界; - 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 链中间动态汇入某个已存在的PCollection | input.apply(Flatten.with(pc2)) |
在 transform 链中间直接汇入Create/Read等根变换 | input.apply(Flatten.with(Create.of(...)))或input.apply(Flatten.with(TextIO.read().from(path))) |
| 需要合并窗口/触发器兼容的多个流式或批量数据源 | 先统一窗口,再使用上述任一方式合并 |
使用FlattenWith时请记住三条核心约束:
- 元素类型必须一致:所有参与合并的
PCollection必须存储相同的数据类型(可使用不同 Coder,但元素类型要一致); - 窗口策略必须兼容:WindowFn 与 Trigger 不兼容会在构建期触发
IllegalStateException; - 合并结果无序:输出
PCollection只是逻辑上的集合,元素顺序不受保证,测试请使用containsInAnyOrder断言。
如需进一步验证行为,可以直接运行 FlattenTest.java 中的testFlattenWithPCollection与testFlattenWithPTransform两个用例,它们是FlattenWith两种重载最权威的行为参照。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Python SDK 实战:使用 FlattenWith 变换合并 PCollection 并优化链式操作
Apache Beam Python SDK 实战:使用 FlattenWith 变换合并 PCollection 并优化链式操作 FlattenWith 是
大数据批处理流处理数据工程Apache Beam Java Kata 实战:用 Flatten 合并多个 PCollection
Apache Beam Java Kata 实战:用 Flatten 合并多个 PCollection Apache Beam 的 Flatten 是合并同类型
Apache Beam Python Katas 实战:用 Flatten 变换合并多个 PCollection
Apache Beam Python Katas 实战:用 Flatten 变换合并多个 PCollection 导读 Flatten 是 Apache Bea
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考