在 Google Cloud Dataflow 上运行 Apache Beam 管道:DataflowRunner 完整实战指南

发布时间:2026/10/6 7:56:18
在 Google Cloud Dataflow 上运行 Apache Beam 管道:DataflowRunner 完整实战指南 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 提供统一编程模型描述批处理与流处理管道而 Runner 负责把同一份管道代码翻译并调度到不同分布式引擎上执行。本篇技术指南以仓库内 53_dataflow_runner.md 为骨架系统讲解如何用 Cloud Dataflow Runner 在 Google Cloud Dataflow 托管服务上运行 Beam 管道从云项目与存储桶准备、Java 依赖声明、Pipeline Options 配置到 Java/Python 双 SDK 的实际运行命令与作业监控并深入仓库源码揭示 Runner 的选项解析与执行机制。读完本文你将掌握一条从零到一、可复制的 Dataflow 作业提交完整流程。一、Cloud Dataflow Runner 概览托管批流一体执行Cloud Dataflow Runner 是为在 Google Cloud Dataflow 服务上运行管道而设计的 Beam Runner。Cloud Dataflow 提供全托管的统一流式与批处理数据处理能力核心特性包括动态工作再平衡dynamic work rebalancing自动感知各 worker 的处理进度把热点数据段的剩余工作迁移给空闲 worker缓解数据倾斜导致的执行拖尾内置自动扩缩容autoscaling根据积压数据量自动调整 worker 数量降低人工运维成本。当你在 Cloud Dataflow 上执行管道时Runner 的职责是将你的代码与全部依赖上传到 Cloud Storage 存储桶staging通过 Dataflow API创建一个 Dataflow 作业job作业在 Google Cloud Platform 的托管计算资源上执行你的管道图。这一点在仓库源码中可以得到印证DataflowRunner.java 中的public class DataflowRunner extends PipelineRunnerDataflowPipelineJob将运行结果建模为DataflowPipelineJob——即作业的提交、轮询与状态管理都由 Runner 与 Dataflow 服务端协作完成。Runner 通过fromOptions(PipelineOptions options)同文件第 364 行解析选项并持有DataflowPipelineOptions options字段第 227 行整个提交链路以选项解析 → 管道翻译 → 作业创建为骨架。二、前置准备配置云项目与资源在执行任何 Dataflow 作业之前需要按官方 Cloud Dataflow 快速入门中 Before You Begin 的步骤完成环境准备选择或创建 Google Cloud Platform Console 项目所有 Dataflow 作业都隶属于某个项目用于计量与权限隔离为项目启用结算billingDataflow 按计算资源与存储计费未启用结算无法创建作业启用所需 Google Cloud API至少包括 Cloud Dataflow、Compute Engine、Stackdriver LoggingCloud Logging、Cloud Storage、Cloud Storage JSON 与 Cloud Resource Manager若管道代码使用了 BigQuery、Pub/Sub 等额外服务需按需再启用对应 API完成 Google Cloud Platform 身份认证本地通过 Application Default Credentials 或gcloud auth application-default login建立凭证使 Runner 有权限创建作业、读写 GCS安装 Google Cloud SDK用于命令行认证、配置与后续作业管理创建 Cloud Storage 存储桶作为代码包、临时文件与输出结果的落盘位置。注意远程执行模式下--project等关键选项不允许为空这正是源码中多处注释Remote execution must check that this option is not None所强调的硬性约束参见 pipeline_options.py。三、Java SDK在 pom.xml 中声明 Runner 依赖使用 Apache Beam Java SDK 时需要在 Java 项目的pom.xml中声明对 Cloud Dataflow Runner 的依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-google-cloud-dataflow-java/artifactId version2.54.0/version scoperuntime/scope /dependency实践要点版本匹配version应与项目使用的 Beam SDK 版本保持一致。原文档示例使用 2.54.0当前仓库开发版本号为2.78.0-SNAPSHOT见 gradle.properties属于未发布的开发快照日常使用请替换为与你 SDK 对应的已发布正式版本号自包含应用确保包含运行所需的全部依赖使应用自包含自执行 JARself-executing JAR某些启动场景如使用调度器触发管道需要把作业打成自执行 JAR。此时需在pom.xml的build项目配置段显式添加生成可执行 JAR 的插件依赖例如maven-shade-plugin并声明Main-Class再配合--runnerDataflowRunner提交。更详细的打包说明可参考官方文档中 Cloud Dataflow Runner 的 Self-executing JAR 一节。四、配置 Pipeline Options控制作业执行细节通过 Pipeline Options 可以控制 Cloud Dataflow 执行作业的方方面面——例如指定管道运行在worker 虚拟机、Cloud Dataflow 服务后端还是本地。两个核心配置入口分别是JavaDataflowPipelineOptions接口org.apache.beam.runners.dataflow.options包PythonGoogleCloudOptions类。4.1 Java 侧DataflowPipelineOptions该接口定义在 DataflowPipelineOptions.java同时继承PipelineOptions、GcpOptions、StreamingOptions、BigQueryOptions、PubsubOptions等一组子接口因而聚合了项目、GCS、流式与各 I/O 相关选项。其核心成员包括选项说明默认行为project云项目 ID云上运行必需Validation.Required通过DefaultProjectFactory自动推断region创建 Dataflow 作业的 Compute Engine 区域由DefaultGcpRegionFactory提供默认区域stagingLocation本地文件代码与依赖的 GCS 暂存路径须以gs://开头未设置时回退到gcpTempLocation并追加/staging后缀见源码第 210-243 行的StagingLocationFactorygcpTempLocationGCS 临时文件路径由GcpOptions提供同时是stagingLocation的回退基础update是否用同名新管道替换正在运行的管道并保留状态布尔开关createFromSnapshot从指定快照创建作业配合快照恢复场景templateLocation生成模板文件的路径本地或 GCSClassic/Flex Template 必需无serviceAccount以指定服务账号运行作业替代默认 GCE robot无labels应用于作业计费记录的标签键值对无dataflowServiceOptions由用户设置的服务端选项使服务端特性可用性脱离 Beam 发布周期无flexRSGoal弹性资源调度目标UNSPECIFIED/SPEED_OPTIMIZED追求更短执行时间/COST_OPTIMIZED追求更低成本UNSPECIFIEDjdkAddOpenModules针对 JDK 16 强封装导致的反射InaccessibleObjectException按module/packagetarget-module格式开放模块无streaming是否为流式管道继承自StreamingOptions布尔开关其中stagingLocation与gcpTempLocation的联动关系在源码中有清晰体现StagingLocationFactory.create先尝试读取gcpTempLocation若缺失或不是合法 GCS 路径会抛出带明确提示的IllegalArgumentExceptionDataflowPipelineOptions.java。因此至少要显式提供一个命令行中常用--gcpTempLocationgs://BUCKET/temp/。4.2 Python 侧GoogleCloudOptionsPython SDK 中pipeline_options.py 第 1014 行定义了class GoogleCloudOptions(PipelineOptions)通过_add_argparse_args暴露命令行参数--dataflow_endpointDataflow API 地址默认https://dataflow.googleapis.com--project拥有该 Dataflow 作业的云项目名远程执行为必填--job_nameCloud Dataflow 作业名称--staging_locationGCS 上暂存代码包的位置未设置时回退到temp_location--temp_locationGCS 上保存临时工作文件的位置远程执行为必填--region创建作业的 Compute Engine 区域。同时在类级定义了OAUTH_SCOPES默认申请 BigQuery、Cloud Platform、Cloud Storage 全控制、用户信息、Datastore、Spanner 等 API 的访问范围pipeline_options.py这也是 Python 管道能直接读写多种 GCP 服务而无需额外配置的前提。五、运行你的管道Java 与 Python 双示例原文档提供了取自 Cloud Dataflow Java/Python 快速入门的 WordCount 示例。执行时占位符PROJECT_ID、BUCKET_NAME、REGION以及 Python 例中的STORAGE_BUCKET、DATAFLOW_REGION需替换为你自己的项目与区域信息。5.1 Java SDKMaven 命令行在word-count-beam目录下执行mvn -Pdataflow-runner compile exec:java \ -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Dexec.args--projectPROJECT_ID \ --gcpTempLocationgs://BUCKET_NAME/temp/ \ --outputgs://BUCKET_NAME/output \ --runnerDataflowRunner \ --regionREGION参数含义分解-Pdataflow-runner激活 Maven profile把 Dataflow Runner 依赖引入运行时 classpath-Dexec.mainClass指定主类这里是 Beam 的org.apache.beam.examples.WordCount--runnerDataflowRunner显式指定 Runner。Java SDK 下该值写作DataflowRunner大小写敏感--project云项目 IDDataflowPipelineOptions 中Validation.Required--gcpTempLocationGCS 临时目录Runner 也会据此推导 staging 位置--outputWordCount 结果输出路径指向gs://BUCKET_NAME/output目录--region作业创建区域建议与数据所在区域就近以降低延迟与成本。5.2 Python SDKpython -m 方式python -m apache_beam.examples.wordcount \ --region DATAFLOW_REGION \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://STORAGE_BUCKET/results/outputs \ --runner DataflowRunner \ --project PROJECT_ID \ --temp_location gs://STORAGE_BUCKET/tmp/参数含义分解--runner DataflowRunnerPython 侧也可以写作dataflow两者等价Runner 别名由PipelineOptions的 runner 映射机制解析--input使用官方公共示例数据gs://dataflow-samples/shakespeare/kinglear.txt《李尔王》全文--output结果写入你自己的存储桶gs://STORAGE_BUCKET/results/outputs--temp_locationPython 远程执行必填用于存放暂存文件--regionDataflow 作业区域。5.3 运行时要点提交后命令即可退出作业在服务端异步执行日志中会打印作业 ID 与 Dataflow Console 链接Java 与 Python 使用的 SDK 版本应分别与各自 Runner 依赖版本匹配否则可能触发协议或 API 不兼容告警若使用调度器定时触发Java 侧需按第三节说明准备自执行 JAR。六、监控 Dataflow 作业作业提交后可通过以下两种官方界面跟踪进度、查看执行细节并接收结果更新Dataflow Monitoring InterfaceWeb 控制台在 Google Cloud Console 的 Dataflow 页面按项目与区域查看作业列表进入单个作业后可以观察管道图各步骤的执行状态与耗时吞吐量与积压backlog指标配合自动扩缩容观察 worker 数量变化worker 日志对接 Stackdriver/Cloud Logging与错误堆栈动态工作再平衡的调度细节。Dataflow Command-line Interfacegcloud CLI使用gcloud dataflow jobs list、gcloud dataflow jobs describe JOB_ID等命令在终端内查询作业状态与元数据适合脚本化监控与 CI 集成。此外Java Runner 返回的DataflowPipelineJob对象本身也封装了作业状态轮询能力见 DataflowRunner.java 泛型签名PipelineRunnerDataflowPipelineJob可在程序内等待job.waitUntilFinish()获取最终状态。七、从源码看 Runner 的执行机制结合仓库源码可以更深入理解选项 → 翻译 → 提交的完整链路选项聚合DataflowRunner.fromOptions(PipelineOptions options)DataflowRunner.java接收PipelineOptions内部通过options.as(DataflowPipelineOptions.class)获得完整配置视图管道翻译Runner 构造DataflowPipelineTranslator.fromOptions(options)同文件第 589 行把 Beam 管道图翻译为 Dataflow 服务可识别的作业描述资源暂存代码与依赖按stagingLocation或推导的gcpTempLocation/staging上传 GCS作业创建与执行调用 Dataflow API 创建作业作业随后在托管 worker 上执行DataflowPipelineJob负责后续状态跟踪。而 Python 侧GoogleCloudOptionspipeline_options.py的_add_argparse_args直接定义了--project、--job_name、--staging_location、--temp_location、--region、--dataflow_endpoint等参数的解析规则并默认申请一组 GCP API 访问 scope——这解释了为何 Python 命令行只需给出上述参数即可完成远程提交。八、能力边界与进一步参考Cloud Dataflow Runner 对 Beam 模型特性的支持程度如状态与定时器、侧输入、窗口与触发器等在不同语言 SDK 上的完备性由Beam Capability Matrix统一维护建议在选型或移植管道前对照该矩阵确认所需能力已受支持。更多运行细节如 Flex Template、自执行 JAR、服务端选项可继续查阅仓库内相关文档Runner 核心源码runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.javaJava 选项接口定义DataflowPipelineOptions.javaPython 选项类定义sdks/python/apache_beam/options/pipeline_options.py仓库当前版本信息gradle.properties总结把 Beam 管道跑上 Cloud Dataflow本质上是配好环境 → 声明依赖 → 配置选项 → 提交作业 → 监控结果五步闭环。本文以原文档为骨架补齐了 Java 的DataflowPipelineOptions关键参数语义、Python 的GoogleCloudOptions参数映射、stagingLocation与gcpTempLocation的回退机制以及 Runner 源码层的提交链路帮助你在托管的批流一体引擎上稳定、可控地运行 Beam 作业。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐在 Google Cloud Dataflow 上运行 Apache Beam 模板python-docs-samples run_template 实战指南在 Google Cloud Dataflow 上运行 Apache Beam 模板python docs samples run_template 实战指南示例工程在 Google Cloud Dataflow 上运行 Apache Beam Python 流水线python-docs-samples 全流程实战指南在 Google Cloud Dataflow 上运行 Apache Beam Python 流水线python docs samples 全流程实战指南 A示例工程Google Cloud Dataflow 实战指南基于 Apache Beam 的流批一体数据管道服务Google Cloud Dataflow 实战指南基于 Apache Beam 的流批一体数据管道服务 Google Cloud Dataflow 是 Go云原生CI/CD运维上一篇微信聊天记录导出教程5 步把微信聊天完整存进电脑下一篇SIG Windows 章程深度解析Kubernetes Windows 节点支持的治理范围与组织架构创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考