- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Partition 是 Apache Beam 中一个非常实用的核心变换(Core Transform):它把类型相同的单个PCollection,按照你提供的分区函数(partitioning function)拆分为固定数量的多个子集合。本文以 Apache Beam 仓库中 Partition Kata 任务 为骨架,完整讲解 Partition 的概念、Kata 的解答实现、底层源码原理与测试验证方式,读完你既能直接完成这道 Kata,也能在真实管道中熟练运用 Partition 做多路分流。
一、Partition 是什么:一个 PCollection 拆成 N 个
在 Beam 编程模型中,PCollection是无界或有界的数据集合。当一批数据需要按某种规则分流到不同处理分支时,可以使用Partition变换。根据 task.md 的定义:
- Partition 适用于存储相同数据类型的
PCollection对象; - 它把一个
PCollection拆分成固定数量(fixed number)的若干较小集合; - 拆分依据是你提供的分区函数——该函数包含决定输入
PCollection元素如何分配到各个结果分区PCollection的逻辑。
典型应用场景包括:按分数段把学生分成几组、按地区把订单分流、按数值范围把日志分级处理等。与GroupByKey(按 Key 聚合)不同,Partition 不做聚合,只是"分类分流";与ParDo多输出(TupleTag)也不同,Partition 无需预先声明多个带标签的输出,而是用"分区索引"直接定位输出集合。
二、Kata 任务要求
本任务位于 learning/katas/java/Core Transforms/Partition/Partition/ 目录,属于学习 Katas 的 "Core Transforms / Partition" 课程。任务内容为:
实现一个
Partition变换,把一个数字PCollection拆分为两个PCollection:第一个包含大于 100的数字,第二个包含其余数字。
从 task-info.yaml 可以看到,这是一个占位符练习(placeholder),Task.java中有一段TODO()等待你补全,完成后由隐藏的单元测试 TaskTest.java 自动校验。
三、Kata 完整解答:Task.java 逐步拆解
完整实现位于 Task.java。我们先看整体骨架:
package org.apache.beam.learning.katas.coretransforms.partition; import org.apache.beam.learning.katas.util.Log; 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.Partition; import org.apache.beam.sdk.transforms.Partition.PartitionFn; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.values.PCollectionList; public class Task { public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline = Pipeline.create(options); PCollection<Integer> numbers = pipeline.apply( Create.of(1, 2, 3, 4, 5, 100, 110, 150, 250) ); PCollectionList<Integer> partition = applyTransform(numbers); partition.get(0).apply(Log.ofElements("Number > 100: ")); partition.get(1).apply(Log.ofElements("Number <= 100: ")); pipeline.run(); } static PCollectionList<Integer> applyTransform(PCollection<Integer> input) { return input .apply(Partition.of(2, (PartitionFn<Integer>) (number, numPartitions) -> { if (number > 100) { return 0; } else { return 1; } })); } }3.1 关键点一:Partition.of(numPartitions, partitionFn)工厂方法
Partition.of(2, (PartitionFn<Integer>) (number, numPartitions) -> { ... })- 第一个参数
2:分区总数(numPartitions),即要把输入拆成几个集合; - 第二个参数:分区函数(partitionFn),对每个元素返回一个分区索引,索引范围必须是
[0, numPartitions-1],即本例中的0或1。
分区逻辑用 Lambda 表达:number > 100返回0(进入第 0 个分区),否则返回1(进入第 1 个分区)。注意 Lambda 需要显式转型为PartitionFn<Integer>。
3.2 关键点二:返回值是PCollectionList<T>
applyTransform的返回类型是PCollectionList<Integer>。PCollectionList是一个按索引访问的 PCollection 集合,通过partition.get(0)和partition.get(1)即可拿到两个子集合,分别输出日志:
partition.get(0).apply(Log.ofElements("Number > 100: ")); partition.get(1).apply(Log.ofElements("Number <= 100: "));这里的Log.ofElements(prefix)是 Katas 提供的日志辅助变换(位于 Log.java),本质是一个PTransform,内部用ParDo+DoFn把元素(可带前缀、可附加窗口信息)打印到日志。
3.3 关键点三:输入数据与预期分流
输入为Create.of(1, 2, 3, 4, 5, 100, 110, 150, 250),共 9 个整数。按照> 100的规则:
- 分区 0(
Number > 100):110, 150, 250 - 分区 1(
Number <= 100):1, 2, 3, 4, 5, 100
注意100本身不满足> 100,因此落入分区 1,这正是边界条件的考察点。
四、底层原理:从源码看 Partition 是如何工作的
Partition 的官方实现位于 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Partition.java。理解它有助于写出正确、健壮的代码。
4.1 类型签名
public class Partition<T> extends PTransform<PCollection<T>, PCollectionList<T>>也就是说,Partition 的输入是PCollection<T>,输出是打包了 N 个PCollection<T>的PCollectionList<T>,元素类型全程保持一致。
4.2 分区函数接口PartitionFn
public interface PartitionFn<T> extends Serializable { int partitionFor(T elem, int numPartitions); }partitionFor接收当前元素与分区总数,返回目标分区的索引(范围[0..numPartitions-1])。由于它继承Serializable,可以安全地序列化到远程 worker 上执行——这是 Beam 分布式执行的基本前提。
此外,源码还提供了带侧输入(side input)的变体PartitionWithSideInputsFn<T>,签名多一个Contextful.Fn.Context c参数,配合Requirements.requiresSideInputs(...)使用,可以在分区决策时参考其他PCollectionView的数据(如阈值)。Kata 用不到,但真实业务中"阈值由配置动态决定"时很有用。
4.3numPartitions的合法性约束
Partition.of(...)构造时会对参数做校验,源码中PartitionDoFn的构造函数明确抛出:
if (numPartitions <= 0) { throw new IllegalArgumentException("numPartitions must be > 0"); }因此分区数必须是正整数,传入0或负数会在管道构造阶段直接失败。
4.4 底层实现:ParDo + TupleTag 多输出
Partition 并不是什么神秘机制,其expand方法本质上是把一个ParDo包装成了多输出形式:
PCollectionTuple outputs = in.apply( ParDo.of(partitionDoFn) .withOutputTags(new TupleTag<Void>() {}, outputTags) .withSideInputs(partitionDoFn.getSideInputs()));- 构造时按分区数生成 N 个
TupleTag(TupleTagList),每个分区对应一个输出标签; - 处理每个元素时,调用分区函数得到索引,再把元素输出到对应标签的集合中(
c.output(typedTag, input)); - 最后把
PCollectionTuple转成PCollectionList,并用输入集合的 Coder统一设置每个输出集合的 Coder。
4.5 分区索引越界的后果
在PartitionDoFn.processElement中,如果分区函数返回的索引不在[0, numPartitions)范围内,会直接抛出IndexOutOfBoundsException:
throw new IndexOutOfBoundsException( "Partition function returned out of bounds index: " + partition + " not in [0.." + numPartitions + ")");也就是说,分区函数必须对每一个元素都返回合法索引,这是编写分区逻辑时最容易出错的地方(例如漏掉某个分支导致返回负数或大于等于 numPartitions 的值)。
4.6 语义保证:Coder、时间戳与窗口
从源码注释可以确认 Partition 的语义保证:
- Coder:默认情况下,输出
PCollectionList中每个集合的 Coder 与输入PCollection相同(pcs.and(outputs.get(typedOutputTag).setCoder(coder))); - 时间戳与窗口:每个输出元素与对应输入元素拥有相同的时间戳,并处于相同的窗口;
- WindowFn:每个输出
PCollection关联的WindowFn与输入一致。
因此 Partition 只做"分流",不改变元素的窗口归属与时间语义,非常适合在窗口化流式管道中做分类处理。
五、单元测试:用 PAssert 验证分区结果
Katas 的隐藏测试 TaskTest.java 展示了标准的分区验证写法:
@Rule public final transient TestPipeline testPipeline = TestPipeline.create(); @Test public void groupByKey() { PCollection<Integer> numbers = testPipeline.apply( Create.of(1, 2, 3, 4, 5, 100, 110, 150, 250) ); PCollectionList<Integer> results = Task.applyTransform(numbers); PAssert.that(results.get(0)) .containsInAnyOrder(110, 150, 250); PAssert.that(results.get(1)) .containsInAnyOrder(1, 2, 3, 4, 5, 100); testPipeline.run().waitUntilFinish(); }要点:
- 使用
TestPipeline(@Rule)驱动测试管道; - 用
PAssert.that(...).containsInAnyOrder(...)断言每个分区的元素无序但完整地等于预期集合——containsInAnyOrder不关心元素顺序,只关心集合内容一致; - 测试输入刻意包含边界值
100,用于检验"大于 100"与"小于等于 100"的边界划分是否正确。
六、如何运行与完成练习
Katas 是为 IntelliJ Education(或 IntelliJ + EduTools 插件)设计的交互式课程,具体配置步骤见 learning/katas/java/README.md:
- 在 IntelliJ Education 中选择Open,打开
learning/katas/java目录; - 按提示Import Gradle project并完成 Gradle 配置;
- 等待 Gradle 构建完成后,在 "Project Structure" 中设置项目 SDK(如 JDK 8);
- 打开 "Project" 工具窗口,切换到Course视图,即可看到 Partition 等课程任务;
- 在
Task.java中替换TODO()占位符,运行测试(隐藏的TaskTest)验证你的实现。
也可以不依赖 IDE,直接运行Task.main(它使用PipelineOptionsFactory.fromArgs(args).create()创建管道并用 Direct Runner 执行),观察控制台日志输出两个分区的元素。
七、举一反三:Partition 的更多用法
掌握 Kata 后,可以把 Partition 推广到更复杂的场景:
按百分比分桶(源码 Javadoc 示例):
PCollectionList<Student> studentsByPercentile = students.apply(Partition.of(10, new PartitionFn<Student>() { public int partitionFor(Student student, int numPartitions) { return student.getPercentile() * numPartitions / 100; // 0..99 } }));基于侧输入动态阈值(PartitionWithSideInputsFn):
PCollectionView<Integer> gradesView = pipeline.apply("grades", Create.of(50)).apply(View.asSingleton()); PCollectionList<Integer> studentsByGrades = pipeline.apply(studentsPercentage) .apply(Partition.of(2, ((elem, numPartitions, ctx) -> { Integer grades = ctx.sideInput(gradesView); return elem < grades ? 0 : 1; }), Requirements.requiresSideInputs(gradesView)));八、总结
- Partition 的定位:把一个同类型
PCollection按自定义分区函数拆成固定数量的子集合,返回PCollectionList<T>; - 核心 API:
Partition.of(numPartitions, partitionFn),分区函数返回[0, numPartitions-1]的索引,numPartitions必须大于 0; - 底层机制:基于
ParDo+TupleTag多输出实现,输出集合沿用输入 Coder、时间戳与窗口语义,索引越界会抛IndexOutOfBoundsException; - 验证方式:
TestPipeline+PAssert.containsInAnyOrder逐分区断言; - 实战价值:Kata 解答(
number > 100 ? 0 : 1)即是最小可运行的分区示例,稍加扩展即可用于分桶、分流、动态阈值等真实场景。
配套练习与源码:任务文档在 task.md,解答在 Task.java,官方变换实现在 Partition.java。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
douyin-downloader 抖音无水印批量下载:从 Cookie 配置到首次入库的上手指南
douyin downloader 抖音无水印批量下载:从 Cookie 配置到首次入库的上手指南 douyin downloader 是一个 Python 开
大数据批处理流处理数据工程小爱音箱接入大模型:MiGPT 部署与配置完整指南
小爱音箱接入大模型:MiGPT 部署与配置完整指南 晚上问小爱同学"为什么天空是蓝色的",它还是那句模板式的客服腔。MiGPT 是一个把小爱音箱接入 ChatG
大数据批处理流处理数据工程使用tradingview-mcp必须知道的4条localhost安全实践
使用tradingview mcp必须知道的4条localhost安全实践 tradingview mcp 是一个把 Claude Code 连接到你本地 Tr
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考