
【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载本指南基于 learning/tour-of-beam/learning-content/introduction/introduction-concepts/runner-concepts/description.md 整理而成。Apache Beam 提供一套可移植的 API 层用于构建复杂的数据并行处理管线Pipeline并允许同一份管线代码运行在多种执行引擎即 Runner之上。本文将以该文档为主线结合本仓库中各 Runner 的源码实现系统讲解 Beam Runner 的核心概念、Direct Runner 的模型校验机制以及 Google Cloud Dataflow、Apache Flink、Apache Spark、Apache Samza、Apache Nemo、Hazelcast Jet 等 Runner 的适用场景与运行方式帮助读者完成 Runner 选型并掌握在 Java、Python、Go 三种 SDK 下的实际运行命令。1. 核心概念Beam Runner 是什么Apache Beam 提供的可移植 API 层允许开发者只编写一次管线代码就将其执行在不同的执行引擎Runner上。这一层的核心概念基于 Beam Model即早先所说的 Dataflow Model而每个 Runner 只是对该模型不同程度的实现。也就是说管线Pipeline描述数据处理的完整逻辑与具体执行环境无关Runner负责把这条管线翻译并调度到某个具体执行引擎上例如本地 JVM、Google Cloud 托管服务、Flink 集群或 Spark 集群不同 Runner 对 Beam Model 的支持程度不同因此同一份代码在不同 Runner 上的行为可能有细微差别这也正是 Direct Runner 存在的重要原因。从源码结构看本仓库在runners/目录下按执行引擎组织各个 Runner 模块例如 runners/direct-java、runners/google-cloud-dataflow-java、runners/flink、runners/spark、runners/samza、runners/jet 等每个模块内部都实现了 Beam 的PipelineRunner接口。例如 FlinkRunner.java 声明为public class FlinkRunner extends PipelineRunnerPipelineResultSparkRunner.java 声明为public final class SparkRunner extends PipelineRunnerSparkPipelineResult可见所有 Runner 都是 PipelineRunner 的实现这一模型在代码层面是统一成立的。2. Direct Runner在本地校验 Beam 模型的正确性2.1 设计目标验证语义而非追求性能Direct Runner 在开发者本机执行管线它的设计目标不是高效执行而是尽可能严格地验证管线是否遵守 Apache Beam 模型防止用户依赖那些模型并不保证的语义。从 runners/direct-java 模块的 DirectRunner.java 源码注释可以确认The DirectRunner is suitable for running a Pipeline on small scale, example, and test data, and should be used for ensuring that processing logic is correct.也就是说Direct Runner 适合小规模、示例与测试数据其核心价值在于保证处理逻辑正确并为后续在分布式后端大规模执行扫清隐患。2.2 模型校验的具体检查项原文档指出 Direct Runner 会执行以下额外检查强制元素不可变性immutability任何 Transform 都不允许在任意时刻修改输入元素也不允许在输出元素被产出之后修改它们强制元素可编码性encodability每个 PCollection 中的全部元素都必须能被该 PCollection 的 Coder 编码与解码任意顺序处理所有环节中的元素都以任意顺序被处理用户函数可序列化DoFn、CombineFn 等用户函数必须能够被序列化。上述检查在源码中有明确对应。在 DirectOptions.java 中isEnforceImmutability()与isEnforceEncodability()两个选项默认均为true分别控制是否校验元素不被篡改、是否校验元素可被 Coder 编码解码。在 DirectRunner.java 中Enforcement枚举实现了ENCODABILITY与IMMUTABILITY两类强制检查其中不可变性检查会通过ImmutabilityCheckingBundleFactory包裹 Bundle 工厂见bundleFactoryFor方法并且只会作用于包含用户函数UDF的变换源码中限定了READ_TRANSFORM_URN与PAR_DO_TRANSFORM_URN。2.3 为什么本地单元测试如此重要原文档特别强调当管线运行在远程集群上时排查失败运行往往非常困难而在本地用 Direct Runner 对管线代码做单元测试则又快又简单还能使用你熟悉的本地调试工具。用 Direct Runner 做测试与开发有助于保证管线在不同 Beam Runner 之间具有鲁棒性。2.4 各 SDK 中的使用方式Go SDK在 Go SDK 中默认 Runner 就是DirectRunner。可直接安装并运行官方 wordcount 示例$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount $ wordcount --input PATH_TO_INPUT_FILE --output counts该示例的源码位于 sdks/go/examples/wordcount/wordcount.go其注释同样强调可以在本地执行管线也可以通过选择其他 Runner 执行这与原文档的表述完全一致。Java SDK使用 Java 时必须在pom.xml中声明对 Direct Runner 的运行时依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-direct-java/artifactId version2.41.0/version scoperuntime/scope /dependency启动程序时通过args设置 Runner--runnerDirectRunner本仓库的官方示例 WordCount.java 也明确说明通过--runnerYOUR_SELECTED_RUNNER即可切换 Runner使用 DirectRunner 时输出为本地文件。Python SDK在 Python SDK 中默认 Runner 同样是DirectRunner。运行 wordcount 示例python -m apache_beam.examples.wordcount --input YOUR_INPUT_FILE --output counts对应示例源码位于 sdks/python/apache_beam/examples/wordcount.py。其底层实现在 sdks/python/apache_beam/runners/direct/direct_runner.py 中其中SwitchingDirectRunner会在FnApiRunner批处理吞吐高与BundleBasedDirectRunner支持流式执行与某些尚未实现的原语之间自动切换这从侧面印证了 Python Direct Runner 在批/流场景下的覆盖策略。3. Google Cloud Dataflow Runner托管式的云上执行3.1 工作原理Google Cloud Dataflow Runner 使用 Cloud Dataflow 托管服务。运行管线时Runner 会把你的可执行代码与依赖上传到一个 Google Cloud StorageGCS桶并创建 Cloud Dataflow 任务Job由该任务在 Google Cloud Platform 的托管资源上执行管线。从 DataflowRunner.java 的源码可以确认其内部大量操作围绕 GCP 资源展开如上传代码、依赖检测、与 Dataflow API 交互等。Dataflow Runner 与服务适用于大规模、持续运行的作业并提供完全托管的服务fully managed service在整个作业生命周期内自动伸缩 worker 数量autoscaling动态工作再平衡dynamic work rebalancing3.2 各 SDK 运行方式Go SDK$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner dataflow \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latestJava SDK在pom.xml中声明依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-google-cloud-dataflow-java/artifactId version2.42.0/version scoperuntime/scope /dependency然后在 Maven JAR 插件中配置mainClassplugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-jar-plugin/artifactId version${maven-jar-plugin.version}/version configuration archive manifest addClasspathtrue/addClasspath classpathPrefixlib//classpathPrefix mainClassYOUR_MAIN_CLASS_NAME/mainClass /manifest /archive /configuration /plugin命令行运行java -jar target/beam-examples-bundled-1.0.0.jar \ --runnerDataflowRunner \ --projectYOUR_GCP_PROJECT_ID \ --regionGCP_REGION \ --tempLocationgs://YOUR_GCS_BUCKET/temp/Python SDK先安装 GCP 相关的扩展组件pip install apache-beam[gcp]再运行python -m apache_beam.examples.wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://YOUR_GCS_BUCKET/counts \ --runner DataflowRunner \ --project YOUR_GCP_PROJECT \ --region YOUR_GCP_REGION \ --temp_location gs://YOUR_GCS_BUCKET/tmp/4. Apache Flink Runner流式优先的高吞吐执行4.1 特性与适用场景Apache Flink Runner 基于 Apache Flink 执行 Beam 管线既可以采用集群执行模式如 YARN / Kubernetes / Mesos也可以采用本地嵌入式模式便于测试管线。Flink Runner 与 Flink 适合大规模、持续运行的作业并提供流式优先streaming-first的运行时同时支持批处理与流数据处理程序同时支持极高吞吐与低事件延迟的运行时具备 exactly-once 处理保证的容错能力流式程序中天然的反压back-pressure机制自定义内存管理可在内存内与内存外数据处理算法之间高效、稳健地切换与 YARN 以及 Apache Hadoop 生态其他组件的集成4.2 各 SDK 运行方式Go SDK需要引入 flink runner 包并通过--endpoint指定 Runner 所在端点github.com/apache/beam/sdks/v2/go/pkg/beam/runners/flink$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner flink \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latest \ --endpointlocalhost:8081Java SDKPortable 模式从 Beam 2.18.0 起Docker Hub 提供了预构建的 Flink Job Service 镜像Flink 1.10、Flink 1.11、Flink 1.12、Flink 1.13、Flink 1.14启动 JobService 端点docker run --nethost apache/beam_flink1.10_job_server:latest使用 PortableRunner 将管线提交到上述端点job_endpoint设为localhost:8099JobService 的默认地址可选设置environment_type为LOOPBACK。示例import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions options PipelineOptions([ --runnerPortableRunner, --job_endpointlocalhost:8099, --environment_typeLOOPBACK ]) with beam.Pipeline(options) as p: ...Java SDK非 Portable 模式在pom.xml中声明依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-flink-1.14/artifactId version2.42.0/version /dependency命令行运行mvn exec:java -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Pflink-runner \ -Dexec.args--runnerFlinkRunner \ --inputFile/path/to/pom.xml \ --output/path/to/counts \ --flinkMasterflink master url \ --filesToStagetarget/word-count-beam-bundled-0.1.jarPython SDK与上述 Portable 模式完全一致先docker run --nethost apache/beam_flink1.10_job_server:latest启动 JobService再通过PortableRunner、job_endpointlocalhost:8099、environment_typeLOOPBACK提交管线。5. Apache Spark Runner与 Spark 生态深度集成5.1 特性Apache Spark Runner 基于 Apache Spark 执行 Beam 管线可以像原生 Spark 应用一样运行以自包含应用部署到本地模式、Spark Standalone 资源管理器或运行在 YARN / Mesos 之上。Spark Runner 提供批处理、流式处理以及批流一体的管线与 RDD 和 DStream 相同的容错保证Spark 提供的安全特性基于 Spark 指标系统的内建指标上报同时上报 Beam Aggregators通过 Spark 的 Broadcast 变量原生支持 Beam side-inputs5.2 各 SDK 运行方式Go SDK需要引入 spark runner 包并通过--endpoint指定端点github.com/apache/beam/sdks/v2/go/pkg/beam/runners/spark$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner spark \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latest \ --endpointlocalhost:8081Java SDK非 Portable 模式启动 JobService 端点二选一使用 Docker推荐docker run --nethost apache/beam_spark_job_server:latest或从 Beam 源码构建./gradlew :runners:spark:3:job-server:runShadow通过 PortableRunner 提交 Python 管线到上述端点job_endpoint设为localhost:8099environment_type设为LOOPBACKPython 示例代码同上文 Flink 部分。在pom.xml中声明 Spark Runner 依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-spark-3/artifactId version2.42.0/version /dependency并使用 Maven Shade 插件对应用 JAR 做 shading避免打包冲突plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId configuration createDependencyReducedPomfalse/createDependencyReducedPom filters filter artifact*:*/artifact excludes excludeMETA-INF/*.SF/exclude excludeMETA-INF/*.DSA/exclude excludeMETA-INF/*.RSA/exclude /excludes /filter /filters /configuration executions execution phasepackage/phase goals goalshade/goal /goals configuration shadedArtifactAttachedtrue/shadedArtifactAttached shadedClassifierNameshaded/shadedClassifierName transformers transformer implementationorg.apache.maven.plugins.shade.resource.ServicesResourceTransformer/ /transformers /configuration /execution /executions /plugin命令行运行mvn compile exec:java -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Dexec.args--runnerSparkRunner --inputFilepom.xml --outputcounts -Pspark-runnerPython SDKpython -m apache_beam.examples.wordcount --input /path/to/inputfile \ --output /path/to/write/counts \ --runner SparkRunner从源码层面看SparkRunner.java 提供了create()、create(SparkPipelineOptions)与fromOptions(PipelineOptions)等多种构造入口并返回SparkPipelineResult表明 Beam 管线在 Spark 上以原生 Spark 应用的形式被执行与追踪。6. Apache Samza Runner大规模有状态流式作业6.1 特性Apache Samza Runner 在 Samza 应用中执行 Beam 管线可本地运行也可以把应用打包成.tgz部署到 YARN 集群或带 Zookeeper 的 Samza standalone 集群。Samza Runner 与 Samza 适合大规模、有状态的流式作业并提供对本地状态基于 RocksDB 存储的一流支持便于高频流式作业快速访问状态支持状态增量 checkpoint 而非全量快照的容错能力使 Samza 能扩展到状态极大的应用完全异步的处理引擎让远程调用更高效灵活的部署模型可在任何带 Zookeeper 的托管环境中运行应用金丝雀发布canaries、升级与回滚等特性支持以最小停机时间支撑超大规模部署原文档中的 Samza 小节仅面向 Java 与 Go SDK 展开。6.2 各 SDK 运行方式Go SDK需要引入 samza runner 包并通过--endpoint指定端点github.com/apache/beam/sdks/v2/go/pkg/beam/runners/samza$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner samza \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latest \ --endpointlocalhost:8081Java SDK在pom.xml中声明 Samza Runner 及其依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-samza/artifactId version2.42.0/version scoperuntime/scope /dependency !-- Samza dependencies -- dependency groupIdorg.apache.samza/groupId artifactIdsamza-api/artifactId version${samza.version}/version /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-core_2.11/artifactId version${samza.version}/version /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-kafka_2.11/artifactId version${samza.version}/version scoperuntime/scope /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-kv_2.11/artifactId version${samza.version}/version scoperuntime/scope /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-kv-rocksdb_2.11/artifactId version${samza.version}/version scoperuntime/scope /dependency命令行运行$ mvn exec:java -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Psamza-runner \ -Dexec.args--runnerSamzaRunner \ --inputFile/path/to/input \ --output/path/to/counts7. Apache Nemo Runner编译期优化与分布式执行7.1 特性Apache Nemo Runner 基于 Apache Nemo 执行 Beam 管线通过 Nemo 编译器中的多种优化 passoptimization passes对 Beam 管线做优化再在 Nemo 运行时上分布式执行。你也可以把应用作为自包含程序部署到本地模式或使用 YARN、Mesos 等资源管理器运行。Nemo Runner 提供批处理与流式管线容错能力与 YARN 以及 Apache Hadoop 生态其他组件的集成对 Nemo 优化器提供的多种优化的支持7.2 Java 运行方式在pom.xml中声明依赖注意需排除 slf4j 冲突dependency groupIdorg.apache.nemo/groupId artifactIdnemo-compiler-frontend-beam/artifactId version${nemo.version}/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-common/artifactId version${hadoop.version}/version exclusions exclusion groupIdorg.slf4j/groupId artifactIdslf4j-api/artifactId /exclusion exclusion groupIdorg.slf4j/groupId artifactIdslf4j-log4j12/artifactId /exclusion /exclusions /dependency原文档建议自包含应用可能更易于管理并能完整使用 Nemo 提供的功能。只需添加上述依赖并用 Maven Shade 插件对应用 JAR 做 shadingShade 配置与第 5.2 节 Spark 的完全相同可复用。命令行运行$ mvn package -Pnemo-runner java -cp target/word-count-beam-bundled-0.1.jar org.apache.beam.examples.WordCount \ --runnerNemoRunner --inputFilepwd/pom.xml --outputcounts8. Hazelcast Jet Runner实验性的高吞吐内存计算8.1 特性与现状说明Hazelcast Jet Runner 基于 Hazelcast Jet 执行 Beam 管线适合大规模持续作业并提供同时支持批处理有界与流式无界数据集同时支持极高吞吐与低事件延迟的运行时流式程序中天然的反压机制带内存存储的分布式大规模并行处理引擎需要特别注意Jet Runner 目前处于 EXPERIMENTAL实验性状态尚不能使用 Jet 的许多能力Jet 本身具备完整的容错支持但 Jet Runner没有作业失败后必须重启Jet 内部性能极高但 Runner 目前无法匹敌因为 Beam 管线的优化/改造surgery尚未完全实现。8.2 Java 运行方式$ mvn package -P jet-runner java -cp target/word-count-beam-bundled-0.1.jar org.apache.beam.examples.WordCount \ --runnerJetRunner --jetLocalMode3 --inputFilepwd/pom.xml --outputcounts其中--jetLocalMode3用于指定 Jet 的本地运行模式。9. Runner 选型速查与实战建议结合原文档与仓库源码可以把上述 Runner 归纳为三类使用场景Runner场景定位关键运行参数备注Direct Runner本地开发、单元测试、模型语义校验--runnerDirectRunnerJava/ 默认Go、Python默认开启不可变性与可编码性检查Dataflow Runner云上大规模持续作业--runnerDataflowRunner--project/--region/--tempLocation托管服务、自动伸缩、动态再平衡Flink Runner流式优先的高吞吐批流作业--runnerFlinkRunner/ PortableRunner --job_endpoint支持 YARN / Kubernetes / 本地嵌入式Spark Runner与 Spark 生态深度集成的批流作业--runnerSparkRunner-Pspark-runner侧输入走 Broadcast 变量Samza Runner大规模有状态流式作业--runnerSamzaRunner-Psamza-runnerRocksDB 本地状态、增量 checkpointNemo Runner追求编译期优化的分布式执行--runnerNemoRunner-Pnemo-runner需 shade 应用 JARJet Runner实验性高吞吐内存计算--runnerJetRunner --jetLocalMode3实验状态作业失败需重启选型建议本地开发与测试首选 Direct Runner。它能提前发现元素被篡改、Coder 不匹配、依赖非模型保证语义等问题且可以直接使用本地调试工具需要托管服务与弹性伸缩时选择 Dataflow适合云上大规模持续作业流式优先、需要 exactly-once 与反压支持时选择 Flink与 Spark 生态RDD/DStream、Broadcast、Spark 指标深度集成时选择 Spark大规模有状态流式作业选择 SamzaRocksDB 本地状态 增量 checkpoint实验性探索可选择 Nemo 或 Jet但要接受 Jet Runner 当前不支持完整容错、作业失败需重启的现实。需要补充说明的是本文涉及的依赖版本号如beam-runners-direct-java2.41.0、beam-runners-google-cloud-dataflow-java2.42.0 等均来自原文档实际使用时应以你当前构建工具解析到的最新版本为准${samza.version}、${nemo.version}、${hadoop.version}等占位符需在你的构建中显式定义。此外Dataflow 运行还需要预先配置 GCP 项目、GCS 桶等环境资源本仓库的 runners/google-cloud-dataflow-java 与 runners/spark、runners/flink 等目录提供了各 Runner 的完整实现可作为进一步研究底层调度与翻译逻辑的入口。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Runner 概念详解从 DirectRunner 到分布式 Runner 的选择与实战Apache Beam Runner 概念详解从 DirectRunner 到分布式 Runner 的选择与实战 Apache Beam 的核心价值在于一次大数据批处理流处理数据工程Warp 补全引擎的 Basic Parser 架构从 Lex 到类型驱动 Full Parse 的递归下降解析器全解析Warp 补全引擎的 Basic Parser 架构从 Lex 到类型驱动 Full Parse 的递归下降解析器全解析 导读 本文深入剖析 Warp 开源仓批处理流处理大数据Apache Beam Direct Runner 完全指南在本地运行与调试 Beam 管道Apache Beam Direct Runner 完全指南在本地运行与调试 Beam 管道 导读 Apache Beam 提供了统一的批处理与流处理编程模型大数据批处理流处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考