ARTICLE DETAIL

资讯详情

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

Apache Beam 入门指南:统一批流处理模型、多语言 SDK 与 Runner 执行体系

Apache Beam 入门指南:统一批流处理模型、多语言 SDK 与 Runner 执行体系 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 是一个用于定义批处理和流处理数据并行管道的统一编程模型并配套提供多语言 SDK 与分布式执行引擎Runner。本文基于当前仓库根目录的 README.md结合仓库中的实际源码目录完整梳理 Beam 编程模型的四个核心概念PCollection、PTransform、Pipeline、PipelineRunner、仓库内置的 Java / Python / Go 等 SDK 体系以及 DirectRunner、FlinkRunner、SparkRunner 等可用 Runner 的源码位置帮助你在进入具体语言 Quickstart 之前建立对 Beam 整体架构的准确认知。一、什么是 Apache BeamBeam 是定义批与流数据并行处理管道的统一模型a unified model for defining both batch and streaming>public static void main(String[] args) { WordCountOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(WordCountOptions.class); runWordCount(options); }即从命令行参数构造并校验PipelineOptions再组装并执行管道。该示例的 Javadoc 中说明输入文件默认是莎士比亚《李尔王》的公开语料可用--inputFile覆盖输出通过--output指定本地文件或输出前缀。此外仓库learning/目录下还有完整的分语言练习体系katas例如 learning/katas/java/Examples/Word Count 与 learning/katas/python/Examples/适合按Introduction → Core Transforms → Common Transforms → Windowing的路径系统学习。第 3 步理解核心概念。即下文展开的PCollection、PTransform、Pipeline与PipelineRunner四件套。三、Beam 模型的四个核心概念README 将 Beam 编程模型的关键概念归纳为四个下面逐一给出定义并对应到本仓库的源码实现。3.1 PCollection数据集合PCollection表示一个数据集合大小可以是有界bounded或无界unbounded——这是批流一体模型中最根本的抽象批处理管道处理有界 PCollection流处理管道处理无界 PCollection其余逻辑保持一致。在 Java SDK 中PCollection继承自PValueBase并实现PValue接口PValue同时继承PInput与POutput即任何值既是某种转换的输入、又是某种转换的输出sdks/java/core/src/main/java/org/apache/beam/sdk/values/PCollection.javapublic class PCollectionT extends PValueBase implements PValuesdks/java/core/src/main/java/org/apache/beam/sdk/values/PValue.javapublic interface PValue extends POutput, PInput同目录下的KV.java、PCollectionView.java、Row.java等则分别对应键值对、侧输入视图、Schema 化的 Row 等常见数据载体。3.2 PTransform转换计算PTransform表示一个把输入 PCollection 转换成输出 PCollection的计算。Java SDK 的所有内建转换都集中在sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/目录下从源码文件名可以直观看到模型覆盖的转换种类逐元素处理DoFn.java用户逻辑基类、MapElements.java、Filter.java、FlatMapElements.java聚合GroupByKey.java、Combine.java、Count.java、Max.java、Mean.java、Distinct.java数据组装Create.java、Impulse.java、Flatten.java、Keys.java、KvSwap.java用户自定义转换通过继承DoFn实现ProcessElement等方法完成WordCount 中的ExtractWordsFn就是典型例子——它把每一行文本切词后receiver.output(word)输出到下游 PCollection并同时通过Metrics.counter/Metrics.distribution上报管道指标。3.3 Pipeline转换与集合构成的 DAGPipeline管理由PTransform和PCollection组成的、准备好被执行的有向无环图DAG。也就是说管道代码只负责声明计算图此时并不发生任何执行——执行交给 Runner。Java 侧的入口类为sdks/java/core/src/main/java/org/apache/beam/sdk/Pipeline.javasdks/java/core/src/main/java/org/apache/beam/sdk/PipelineResult.java3.4 PipelineRunner指定在哪里、如何执行PipelineRunner决定管道的执行位置与方式。其抽象基类定义在 sdks/java/core/src/main/java/org/apache/beam/sdk/PipelineRunner.javapublic abstract class PipelineRunnerResultT extends PipelineResult每个 Runner 实现把自己的管道翻译为目标后端的执行计划。四、Beam 服务的三类用户README 特别指出Beam 面向三类背景与诉求截然不同的用户这也是理解整个仓库分层的钥匙终端用户End Users用现有 SDK 写管道、用现有 Runner 跑只关心应用逻辑希望其他一切自动生效。对应上文 Quick Start 路径与examples/、learning/资源。SDK 编写者SDK Writers为特定语言社区开发 Beam SDKJava、Python、Go 等是语言极客希望被屏蔽各 Runner 的实现细节。仓库中 sdks/ 即各 SDK 的实现其中 Java SDK 的harness/、expansion-service/等子模块体现了 SDK 与服务化扩展能力的分层。Runner 编写者Runner Writers拥有分布式执行环境希望支持针对 Beam 模型编写的程序希望被屏蔽多 SDK 细节。仓库中 runners/ 即各 Runner 的实现。五、仓库内置的 Runner 清单与源码位置README 列出了当前可用的 PipelineRunner 及其执行后端。结合本仓库runners/目录结构对应关系如下从目录结构看README 所列每个 Runner 都有独立的源码子目录Runner执行后端仓库源码目录DirectRunner本地机器runners/direct-javaPrismRunner本地机器基于 Beam Portability容器化runners/prism/javaDataflowRunnerGoogle Cloud Dataflow 服务runners/google-cloud-dataflow-javaFlinkRunnerApache Flink 集群runners/flinkSparkRunnerApache Spark 集群runners/sparkJetRunnerHazelcast Jet 集群runners/jetTwister2RunnerTwister2 集群runners/twister2其中值得注意的两点DirectRunner适合本地开发与单元测试其 Java 实现入口为 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectRunner.javapublic class DirectRunner extends PipelineRunnerDirectPipelineResult。PrismRunner走的是 Beam Portability可移植性路线管道以容器化的外部进程方式执行本地即可验证跨 SDK/跨 Runner 的 gRPC 协议链路实现位于 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismRunner.java。此外 runners/portability/java 提供了 Portability 相关的测试与工具脚本如 runners/portability/test_pipeline_jar.sh。从 README 的表述看FlinkRunner 与 JetRunner 的代码分别由 dataArtisans/flink-dataflow 和 hazelcast/hazelcast-jet 项目捐赠而来Twister2Runner 来自 DSC-SPIDAL/twister2现已全部纳入 Beam 仓库统一维护。另外runners/目录下还有core-java、kafka-streams、local-java、java-fn-execution等模块从目录结构看分别对应 Runner 公共基础设施、Kafka Streams 执行扩展以及跨语言函数执行fn-execution等支撑组件属于 Runner 体系的内建能力层。六、版本与构建支持矩阵以当前仓库实际内容为准gradle.properties 中声明的关键版本信息为version2.78.0-SNAPSHOT sdk_version2.78.0.dev flink_versions1.19,1.20,2.0,2.1,2.2 spark_versions3,4 python_versions3.10,3.11,3.12,3.13,3.14也就是说本仓库对应 Beam 2.78.0 的开发快照版本构建层面同时维护 Flink 1.19~2.2 与 Spark 3/4 两套 Flink/Spark Runner 变体Python SDK 覆盖 3.10 到 3.14。README 中还附有 Maven Central、PyPI 等构件版本徽标说明 Java/Python/Go SDK 均以标准包管理渠道对外发布。构建与测试 Beam 本身的操作规范见 CONTRIBUTING.md 与 CI.md。七、学习资源与仓库内延伸阅读README 的 Learn More 部分列出了社区维护的学习资源官网文档、Java/Python/Go Quickstart、Tour of Beam 交互式学习、Beam Quest 认证、社区指标站点。落到本仓库内部以下路径可以构成自洽的学习闭环示例管道examples/java含WordCount.java等经典示例、examples/python、examples/goGo 示例位于 sdks/go/examples分章练习kataslearning/katas 提供 Java / Python / Go / Kotlin 四种语言版本的课程章节覆盖 Introduction、Core Transforms、Common Transforms、Windowing、Triggers、IO 等交互 Playgroundplayground 目录包含前端Flutter/Dart、后端Go、gRPC 协议playground/api/v1/api.proto与多语言容器化执行环境可用于在浏览器内运行 Beam 管道AI/ML 扩展examples/notebooks 中的 beam-ml 系列 NotebookRunInference、模型刷新、数据预处理等与 sdks/java/ml 模块模型与协议model/ 目录下的 proto 定义pipeline、job-management、fn-execution、interactive是理解 Beam 模型如何跨进程/跨语言传输的一手资料贡献者文档contributor-docs/ 中的发布指南、依赖升级规范、RC 测试流程等。八、小结回到 README 的核心脉络Apache Beam 统一模型PCollection PTransform Pipeline× 多语言 SDKJava/Python/Go× 多后端 RunnerDirect/Flink/Spark/Dataflow/Jet/Twister2/Prism。三类用户各取所需终端用户从 WordCount 示例与 katas 课程入手即可写管道跑 RunnerSDK 编写者与 Runner 编写者则分别以sdks/与runners/目录下的模块为骨架进行扩展。理解这一声明式 DAG 可插拔执行的分层设计后再去阅读某一具体语言 Quickstart 或某个 Runner 的实现都会事半功倍。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Nova框架本地化攻略多语言支持与文化适配技巧Nova框架本地化攻略多语言支持与文化适配技巧 Nova框架是一款基于Unity的视觉小说VN/文字冒险游戏AVG开发框架为开发者提供了便捷的多语言 datasets与Apache Beam集成批处理与流处理统一 datasets与Apache Beam集成批处理与流处理统一 你是否还在为机器学习数据处理中的批处理与流处理分裂而困扰是否希望用一套工具链解决从历史数据集数据工程机器学习5分钟上手 AGENTS.md一份配置全公司 AI 编码工具通用5分钟上手 AGENTS.md一份配置全公司 AI 编码工具通用 换了编辑器规则就丢了Cursor 用 .cursorrulesAider 用 .ai文档教程AI Agent上一篇在Windows平台上构建libpcap库的完整指南下一篇Daft项目深度解析分布式数据框架的分区机制详解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表