ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Beam 核心变换实战:使用 Partition 将 PCollection 拆分为多个输出集合(Java Kata 详解)

Apache Beam 核心变换实战:使用 Partition 将 PCollection 拆分为多个输出集合(Java Kata 详解) 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载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 可以看到这是一个占位符练习placeholderTask.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); PCollectionInteger numbers pipeline.apply( Create.of(1, 2, 3, 4, 5, 100, 110, 150, 250) ); PCollectionListInteger partition applyTransform(numbers); partition.get(0).apply(Log.ofElements(Number 100: )); partition.get(1).apply(Log.ofElements(Number 100: )); pipeline.run(); } static PCollectionListInteger applyTransform(PCollectionInteger input) { return input .apply(Partition.of(2, (PartitionFnInteger) (number, numPartitions) - { if (number 100) { return 0; } else { return 1; } })); } }3.1 关键点一Partition.of(numPartitions, partitionFn)工厂方法Partition.of(2, (PartitionFnInteger) (number, numPartitions) - { ... })第一个参数2分区总数numPartitions即要把输入拆成几个集合第二个参数分区函数partitionFn对每个元素返回一个分区索引索引范围必须是[0, numPartitions-1]即本例中的0或1。分区逻辑用 Lambda 表达number 100返回0进入第 0 个分区否则返回1进入第 1 个分区。注意 Lambda 需要显式转型为PartitionFnInteger。3.2 关键点二返回值是PCollectionListTapplyTransform的返回类型是PCollectionListInteger。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内部用ParDoDoFn把元素可带前缀、可附加窗口信息打印到日志。3.3 关键点三输入数据与预期分流输入为Create.of(1, 2, 3, 4, 5, 100, 110, 150, 250)共 9 个整数。按照 100的规则分区 0Number 100110, 150, 250分区 1Number 1001, 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 PartitionT extends PTransformPCollectionT, PCollectionListT也就是说Partition 的输入是PCollectionT输出是打包了 N 个PCollectionT的PCollectionListT元素类型全程保持一致。4.2 分区函数接口PartitionFnpublic interface PartitionFnT extends Serializable { int partitionFor(T elem, int numPartitions); }partitionFor接收当前元素与分区总数返回目标分区的索引范围[0..numPartitions-1]。由于它继承Serializable可以安全地序列化到远程 worker 上执行——这是 Beam 分布式执行的基本前提。此外源码还提供了带侧输入side input的变体PartitionWithSideInputsFnT签名多一个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 TupleTagVoid() {}, outputTags) .withSideInputs(partitionDoFn.getSideInputs()));构造时按分区数生成 N 个TupleTagTupleTagList每个分区对应一个输出标签处理每个元素时调用分区函数得到索引再把元素输出到对应标签的集合中c.output(typedTag, input)最后把PCollectionTuple转成PCollectionList并用输入集合的 Coder统一设置每个输出集合的 Coder。4.5 分区索引越界的后果在PartitionDoFn.processElement中如果分区函数返回的索引不在[0, numPartitions)范围内会直接抛出IndexOutOfBoundsExceptionthrow 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() { PCollectionInteger numbers testPipeline.apply( Create.of(1, 2, 3, 4, 5, 100, 110, 150, 250) ); PCollectionListInteger 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(); }要点使用TestPipelineRule驱动测试管道用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 示例PCollectionListStudent studentsByPercentile students.apply(Partition.of(10, new PartitionFnStudent() { public int partitionFor(Student student, int numPartitions) { return student.getPercentile() * numPartitions / 100; // 0..99 } }));基于侧输入动态阈值PartitionWithSideInputsFnPCollectionViewInteger gradesView pipeline.apply(grades, Create.of(50)).apply(View.asSingleton()); PCollectionListInteger 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按自定义分区函数拆成固定数量的子集合返回PCollectionListT核心 APIPartition.of(numPartitions, partitionFn)分区函数返回[0, numPartitions-1]的索引numPartitions必须大于 0底层机制基于ParDoTupleTag多输出实现输出集合沿用输入 Coder、时间戳与窗口语义索引越界会抛IndexOutOfBoundsException验证方式TestPipelinePAssert.containsInAnyOrder逐分区断言实战价值Kata 解答number 100 ? 0 : 1即是最小可运行的分区示例稍加扩展即可用于分桶、分流、动态阈值等真实场景。配套练习与源码任务文档在 task.md解答在 Task.java官方变换实现在 Partition.java。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐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大数据批处理流处理数据工程上一篇dbrx-base-FP8-KVAMD革命性FP8量化大模型4倍内存优化提升推理效率下一篇如何快速掌握asdf-vm构建现代化多语言版本管理平台的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表