批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载本文是 Apache Beam 官方文档《Create Your Pipeline》见仓库 create-your-pipeline.md的深度扩展解读。核心围绕 Beam 管道生命周期五步法展开创建Pipeline对象、用Read/Create引入数据、施加 transforms、输出最终数据、运行管道。读完本文你将掌握用 Beam Java SDK 从零搭建一条可运行管道WordCount 式流程的完整套路并理解每一步背后的源码实现机制。管道构建的核心思想程序即图运行即执行Apache Beam 的编程模型用一个核心抽象贯穿始终你的程序表达一条数据处理管道从起点到终点。这条管道由两类对象构成——PTransform转换和PCollection数据集合二者共同组成一张有向无环图DAG。这一点在 Pipeline.java 的类注释中写得非常明确A {link Pipeline} manages a directed acyclic graph of {link PTransform PTransforms}, and the {link PCollection PCollections} that the {link PTransform PTransforms} consume and produce.构建管道时需要执行的通用步骤如下创建一个Pipeline对象使用Read或Create这类根转换为管道数据创建一个或多个PCollection对每个PCollection施加transforms——转换可以改变、过滤、分组、分析或以其他方式处理PCollection中的元素每次转换都会产出一个新的输出PCollection你可以继续对它施加更多转换直到处理完成写出或以其他方式输出最终转换后的PCollection运行管道。需要特别强调一个关键机制这也是 Beam 与普通函数式代码最大的区别施加转换并不会立即执行任何计算。从源码看Pipeline.applyInternal见 Pipeline.java所做的只是把转换节点推入TransformHierarchy并调用transform.expand(input)完成图上的接线真正的执行要等到run()被调用、runner 拿到整张图之后才开始。因此你可以把构建阶段理解为纯粹地在内存中搭建一张计算图。创建 Pipeline 对象一切的起点一个 Beam 程序通常从创建Pipeline对象开始。在 Beam SDK 中每条管道由一个显式的Pipeline类型对象表示。每个Pipeline对象都是一个独立实体封装了管道所操作的数据以及施加于这些数据之上的转换。要创建管道你需要声明一个Pipeline对象并传入一些配置选项// Start by defining the options for the pipeline. PipelineOptions options PipelineOptionsFactory.create(); // Then create the pipeline. Pipeline p Pipeline.create(options);PipelineOptions 与 PipelineOptionsFactoryPipelineOptions是管道配置的统一载体它决定了 runner 选择、执行环境、资源参数等关键行为。PipelineOptionsFactory.create()会返回一个实现了PipelineOptions接口的代理对象并自动把appName设置为调用方类的简单类名见 PipelineOptionsFactory.java。实际工程中更常用的方式是从命令行解析配置参数PipelineOptionsFactory支持 GNU 风格的参数解析fromArgs见 PipelineOptionsFactory.javapublic static void main(String[] args) { PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); Pipeline p Pipeline.create(options); }命令行参数支持多种形式--projectMyProject简单属性把project设置为MyProject--readOnlytrue布尔属性--readOnly是布尔属性简写等价于置true--x1 --x2 --x3列表式简单属性得到[1, 2, 3]--x1,2,3是它的逗号分隔简写--complexObject{key1:value1,...}复杂类型用 JSON 格式传入。默认启用严格解析参数必须符合--booleanArgName或--argNameargValue的格式空参数会被忽略可以通过Builder.withoutStrictParsing()关闭严格模式以上规则均出自 PipelineOptionsFactory.java 的文档说明。Pipeline.create 背后发生了什么Pipeline.create(options)是静态工厂方法见 Pipeline.java其实现会先调用PipelineRunner.fromOptions(options)根据选项确定将要使用的 runner然后构造一个新的Pipeline。构造函数内部创建了一个TransformHierarchy转换层级后续所有apply调用都会登记到这个层级中最终run()时由 runner 遍历该层级执行。此外Pipeline还内置了CoderRegistry与SchemaRegistry分别通过getCoderRegistry()、getSchemaRegistry()懒加载见 Pipeline.java用于为PCollection推断和注册编码器Coder这解释了为什么很多场景下你不需要手动指定 Coder。把数据读入管道Read 与 Create 两类根转换要创建管道的初始PCollection你需要对管道对象施加一个根转换root transform。根转换可以从外部数据源、也可以从你指定的本地内存数据中创建PCollection。Beam SDK 中有两种根转换Read转换从外部数据源读取数据例如文本文件或数据库表Create转换从内存中的java.util.Collection创建PCollection。下面的例子演示如何施加一个TextIO.Read根转换从文本文件读取数据。转换被施加到Pipeline对象p上返回一个PCollectionString形式的数据集PCollectionString lines p.apply( ReadLines, TextIO.read().from(gs://some/inputData.txt));文本 I/O 的实际行为TextIO.read()的默认配置可以在 TextIO.java 中看到压缩方式默认为Compression.AUTO自动探测 gzip 等压缩格式、skipHeaderLines为 0、EmptyMatchTreatment.DISALLOW匹配不到任何文件时直接报错。TextIO还支持丰富的配置比如PCollectionString lines p.apply(ReadLines, TextIO.read() .from(gs://bucket/path/*.txt) // 支持 glob 通配 .withCompression(Compression.GZIP) .withSkipHeaderLines(1) // 跳过首行表头 .withEmptyMatchTreatment(EmptyMatchTreatment.ALLOW));.from()传入的路径支持通配符如file*.txtBeam 会展开为多个文件分片读取。用 Create 构造内存数据当没有外部文件依赖时——尤其是单元测试场景——Create是最佳选择。官方源码注释明确指出 Create 的适用场景A good use for Create is when a PCollection needs to be created without dependencies on files or other external entities. This is especially useful during testing.见 Create.java。PCollectionString words p.apply(Create.of(hello, world, beam)); // 也支持 Map产出 KV 类型 MapString, Integer map ...; PCollectionKVString, Integer entries p.apply(Create.of(map).withCoder( KvCoder.of(StringUtf8Coder.of(), BigEndianIntegerCoder.of())));需要留意两点其一Create只适合小规模内存数据集源码注释中的 Caveat 明确写出 Create only supports small in-memory datasets见 Create.java其二如果所有元素运行时类型相同且该类型注册了默认 CoderCreate可以自动推断 Coder否则必须显式调用.withCoder(...)指定编码方式见 Create.java。施加转换处理管道数据PTransform 的 apply 语义你可以使用 Beam SDK 提供的各种 transforms 来操纵数据。做法是对每个想要处理的PCollection调用apply方法并把期望的转换对象作为参数传入。下面的代码展示了对一个字符串PCollection施加转换的过程。该转换是用户自定义转换反转每个字符串的内容并输出一个包含反转后字符串的新PCollection。输入是名为words的PCollectionString代码把一个名为ReverseWords的PTransform实例传给apply返回值保存为名为reversedWords的PCollectionStringPCollectionString words ...; PCollectionString reversedWords words.apply(new ReverseWords());自定义 PTransform 的标准写法ReverseWords这类自定义转换通常写成内部类并继承PTransformstatic class ReverseWords extends PTransformPCollectionString, PCollectionString { Override public PCollectionString expand(PCollectionString input) { return input.apply(ParDo.of(new DoFnString, String() { ProcessElement public void processElement(Element String word, OutputReceiverString out) { out.output(new StringBuilder(word).reverse().toString()); } })); } }apply方法在底层统一委托给Pipeline.applyTransform(name, input, transform)见 Pipeline.java它会把转换注册进TransformHierarchy并执行expand完成图的拼接。你给转换起的名字如ReadLines、new ReverseWords()的类名会被用于监控界面、日志以及在管道更新时稳定标识图中的节点见 Pipeline.java。组合出完整的处理链处理链的核心在于每次 apply 产出一个新 PCollection再继续 apply。以官方Pipeline类注释中的示例为蓝本见 Pipeline.java一个完整的词频统计链如下// 多个根转换可以并存 PCollectionString lines p.apply(TextIO.read().from(gs://bucket/dir/file*.txt)); PCollectionString moreLines p.apply(TextIO.read().from(gs://bucket/other/dir/*.txt)); // 合并多个 PCollection PCollectionString allLines PCollectionList.of(lines).and(moreLines) .apply(Flatten.StringpCollections()); // 逐级施加转换 PCollectionKVString, Integer wordCounts allLines .apply(ParDo.of(new ExtractWords())) .apply(Count.StringperElement()); PCollectionString formattedWordCounts wordCounts.apply(ParDo.of(new FormatCounts()));这套根转换 → 中间转换链 → 输出的模式就是 Beam 管道的基本骨架。更复杂的场景分组GroupByKey、窗口Window、触发Trigger也都是以同样的apply语义叠加进来的。写出最终管道数据Write 转换当管道完成所有转换后通常需要输出结果。要对管道最终的PCollection做输出需要对这个PCollection施加一个Write转换。Write转换可以把PCollection的元素输出到外部数据汇data sink例如数据库表。你可以在管道的任何时刻使用Write输出数据不过通常是在管道末尾写出。下面的例子展示了如何施加TextIO.Write转换把一个字符串PCollection写到文本文件PCollectionString filteredWords ...; filteredWords.apply(WriteMyFile, TextIO.write().to(gs://some/outputData.txt));TextIO.write()见 TextIO.java常用配置还包括分片、文件后缀等filteredWords.apply(WriteMyFile, TextIO.write() .to(gs://bucket/output/wordcounts) .withNumShards(3) // 输出分片数 .withSuffix(.txt)); // 输出文件后缀注意TextIO.write().to()指定的是路径前缀实际会按分片策略生成wordcounts-00000-of-00003.txt之类的多个文件。写入完成后文件通常会伴随一个_SUCCESS标记文件供下游判断任务是否完整结束。除文本文件外Beam 还提供 Avro、Parquet、BigQuery、JDBC、Kafka 等大量Write实现接口形态与TextIO一致。运行你的管道run 与 waitUntilFinish管道构建完成后使用run方法执行它。管道是异步执行的你创建的程序把管道规格说明发送给一个管道 runnerpipeline runner由 runner 构建并实际运行一系列管道操作。p.run();如果你需要阻塞式执行可以追加waitUntilFinish方法p.run().waitUntilFinish();run 方法内部的执行链路run()的实现在 Pipeline.java它通过PipelineRunner.fromOptions(options)依据选项选出 runnerDirectRunner、DataflowRunner、FlinkRunner、SparkRunner 等随后依次执行validate(options)校验选项、validateErrorHandlers()最后调用runner.run(this)。如果用户代码抛出异常会被包装成PipelineExecutionException抛出其getCause()可拿到原始异常。waitUntilFinish()定义在 PipelineResult.java会阻塞直到管道到达终态如DONE、FAILED、CANCELLED并返回当前状态。PipelineResult.getState()可用于轮询式地查询执行状态见 PipelineResult.java。runner 的选择runner 由PipelineOptions决定不同 runner 带来完全不同的执行形态runner典型用途DirectRunner本地单机执行用于单元测试与本地开发调试默认 runnerDataflowRunner在 Google Cloud Dataflow 托管服务上运行FlinkRunner/SparkRunner/SamzaRunner等在对应分布式计算引擎集群上运行对于本地验证最简单的运行方式是直接使用默认的 DirectRunner 并配合waitUntilFinish()p.run().waitUntilFinish();对于分布式环境则通过命令行参数切换 runner例如# 使用 DataflowRunner 提交到 Google Cloud java -cp target/word-count-bundled-0.1.jar \ com.example.WordCount \ --runnerDataflowRunner \ --projectmy-project \ --stagingLocationgs://my-bucket/staging \ --inputFilegs://my-bucket/input.txt \ --outputgs://my-bucket/output实战组合一条完整管道的构建与运行把前面五个步骤串起来得到一条可运行的最小完整管道对应官方示例 WordCount.java 的简化版public class ReverseWordsPipeline { static class ReverseWords extends PTransformPCollectionString, PCollectionString { Override public PCollectionString expand(PCollectionString input) { return input.apply(ParDo.of(new DoFnString, String() { ProcessElement public void processElement(Element String word, OutputReceiverString out) { out.output(new StringBuilder(word).reverse().toString()); } })); } } public static void main(String[] args) { // 第 1 步创建 Pipeline 对象从命令行解析配置 PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); Pipeline p Pipeline.create(options); // 第 2 步读取输入得到初始 PCollection PCollectionString lines p.apply( ReadLines, TextIO.read().from(options.as(MyOptions.class).getInputFile())); // 第 3 步施加转换处理数据 PCollectionString reversedWords lines.apply(new ReverseWords()); // 第 4 步写出最终结果 reversedWords.apply(WriteMyFile, TextIO.write().to(options.as(MyOptions.class).getOutput())); // 第 5 步运行管道阻塞等待完成 p.run().waitUntilFinish(); } }自定义选项接口MyOptions的典型定义如下public interface MyOptions extends PipelineOptions { Description(Path to the input file) String getInputFile(); void setInputFile(String value); Description(Path of the output directory) String getOutput(); void setOutput(String value); }这样运行时通过--inputFile... --output...即可注入参数也可以继续叠加--runner...切换执行引擎。测试你的管道构建与验证的闭环管道代码写完并运行之前官方强烈建议先做本地单元测试详见 test-your-pipeline.md。Beam 的间接执行模型用户代码构造管道图、由远程 runner 执行使远程调试成本较高本地测试往往更快更简单。测试的基本模式是创建一个TestPipeline它内部自动处理PipelineOptions替代Pipeline.create用静态已知输入数据 Create转换构造PCollection对输入施加待测转换保存输出PCollection用PAssert校验输出内容符合预期。public class CountTest { static final ListString WORDS Arrays.asList( hi, there, hi, hi, sue, bob, hi, sue, , , ZOW, bob, ); public void testCount() { Pipeline p TestPipeline.create(); // 用 Create 构造输入 PCollection PCollectionString input p.apply(Create.of(WORDS)); // 施加待测转换 PCollectionKVString, Long output input.apply(Count.StringperElement()); // 断言输出内容顺序无关 PAssert.that(output).containsInAnyOrder( KV.of(hi, 4L), KV.of(there, 1L), KV.of(sue, 2L), KV.of(bob, 2L), KV.of(, 3L), KV.of(ZOW, 1L)); p.run(); } }使用PAssert的 Java 代码需要引入 JUnit 与 Hamcrest 依赖在pom.xml中追加org.hamcrest:hamcrest:2.2scope为test。测试通过TestPipeline.create()时若发现管道尚未运行测试框架还会强制要求补上run()避免构建了图却从未执行的无效测试对应 TestPipeline.java 中的PipelineAbandonedNodeEnforcement校验。本地用 DirectRunner 测试通过后再以小规模数据切到目标 runner如本地 Flink 集群做集成验证最后再上生产规模——这是官方推荐的稳妥路径。小结与延伸阅读一条 Beam 管道的生命周期可以浓缩为一句话创建Pipeline→ 用Read/Create引入数据 → 用apply叠加 transforms → 用Write输出 →run()交给 runner 执行。理解构建图与执行图分离这一本质你就掌握了 Beam 的核心心智模型。想继续深入可以阅读仓库中的相关文档Programming Guide创建管道、配置管道选项、施加 transforms 的完整细节Test your pipeline本地测试的完整方法论与PAssert用法Design your pipeline管道图设计的思考框架核心源码Pipeline类Pipeline.java、选项工厂PipelineOptionsFactory.java、Create转换Create.java与文本 I/OTextIO.java可运行的完整示例examples/java 目录下提供 WordCount 等大量开箱即用的管道程序。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam 管道构建全指南从 Pipeline 对象、PCollection 到运行执行的完整流程Apache Beam 管道构建全指南从 Pipeline 对象、PCollection 到运行执行的完整流程 Apache Beam 提供统一的批处理与流处大数据批处理流处理数据工程Apache Beam 流水线Pipeline入门从 DAG 抽象模型到 Python 实战配置Apache Beam 流水线Pipeline入门从 DAG 抽象模型到 Python 实战配置 Apache Beam 的 Pipeline 对象是一个Kedro 数据处理管道实战从 Node 到 Pipeline 的完整构建与运行指南Kedro 数据处理管道实战从 Node 到 Pipeline 的完整构建与运行指南 导读 本文以 Kedro 官方 spaceflights 教程的数据处理数据工程工作流自动化上一篇Darknet跨平台部署终极指南Windows、Linux、macOS兼容性详解下一篇Scrapy框架深度解析Easy-scraping-tutorial企业级爬虫开发指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考