首页
/
行业洞察
/
正文
INDUSTRY INSIGHT · 深度
使用 Google Cloud Dataflow Runner 运行 Apache Beam 管道:配置、选项与最佳实践
📅 2026/10/10 2:43:51
✍️ 爱科研究院
👁 阅读 3,247
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读本文面向希望在 Google Cloud 托管环境中运行 Apache Beam 管道的开发者系统讲解Cloud Dataflow Runner的完整使用流程从环境前置条件、Java/Python 两种 SDK 的依赖配置到自包含可执行 JAR 的打包提交再到runner、project、region、streaming、tempLocation等核心 Pipeline 选项的逐项解析最后覆盖作业监控、阻塞式执行与流式执行的注意事项。读完本文你将掌握在 Cloud Dataflow 服务上提交批处理与流处理作业的完整实战能力并能结合本仓库的 Runner 源码理解底层行为。Cloud Dataflow Runner 是什么Google Cloud Dataflow Runner 是 Apache Beam 的托管式执行后端它使用 Cloud Dataflow 托管服务当你使用 Cloud Dataflow 服务运行管道时Runner 会把你的可执行代码与依赖上传到 Google Cloud Storage 存储桶然后创建一个 Cloud Dataflow 作业由该作业在 Google Cloud Platform 的托管资源上执行你的管道。从源码角度看该 Runner 的核心实现位于 DataflowRunner.java它继承自PipelineRunnerDataflowPipelineJob通过DataflowPipelineTranslator将管道翻译成 Dataflow 作业表示再提交给 Dataflow 服务执行见该类头部注释。其运行入口run(Pipeline pipeline)会依次完成管道优化如ProjectionPushdownOptimizer、实验特性开关设置、作业创建与提交等步骤。Cloud Dataflow Runner 与服务特别适合大规模、持续运行的作业并带来三项核心能力完全托管fully managed的服务无需自行管理计算资源集群自动扩缩容autoscaling在作业整个生命周期内根据负载自动调整 worker 数量动态工作再平衡dynamic work rebalancing自动处理数据倾斜避免出现落后分片导致的拖尾。Beam 官方还维护了一张 Beam Capability Matrix 文档记录了 Cloud Dataflow Runner 支持的能力矩阵在选用 Runner 前可先对照查阅。前置条件与初始化设置要使用 Cloud Dataflow Runner必须首先完成所选语言对应的 Cloud Dataflow quickstart 中Before you begin一节列出的设置包括以下六步选择或创建一个 Google Cloud Platform Console 项目。为该项目启用结算billing。启用所需 Google Cloud APICloud Dataflow、Compute Engine、Stackdriver Logging、Cloud Storage、Cloud Storage JSON 以及 Cloud Resource Manager。如果管道代码中还使用了其他服务如 BigQuery、Cloud Pub/Sub 或 Cloud Datastore则需要额外启用对应 API。完成 Google Cloud Platform 身份认证gcloud auth login或配置服务账号。安装 Google Cloud SDK。创建一个 Cloud Storage 存储桶用于存放临时文件与代码包。这些前置条件在代码层面也有印证Java 侧DataflowRunner.fromOptions()在构造 Runner 时会强制校验appName、region、gcpTempLocation、stagingLocation等必需选项缺少即抛出IllegalArgumentException见 DataflowRunner.java。指定 Runner 依赖Java在 pom.xml 中声明依赖使用 Java SDK 时需要在pom.xml中显式声明对 Cloud Dataflow Runner 的依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-google-cloud-dataflow-java/artifactId version{{ param release_latest }}/version scoperuntime/scope /dependency其中version应替换为你使用的 Beam 版本本仓库当前源码版本见 gradle.properties。scope设为runtime表示该依赖只在运行期需要编译期不需要直接引用其类。注意本节内容不适用于 Beam SDK for Python。Python 用户无需声明此类 Maven 依赖Runner 能力随apache_beam包内置提供。Python无需额外依赖Python 侧没有与 Maven 依赖对应的概念。Runner 通过--runnerDataflowRunner直接使用相关选项定义在 pipeline_options.py 的GoogleCloudOptions类中它集中定义了project、job_name、staging_location、temp_location、region、dataflow_endpoint以及一组 OAuth scope 常量。打包自执行 JARJava在某些场景下例如使用 Apache Airflow 之类的调度器启动管道你需要一个自包含self-contained的可执行应用。为此可以在pom.xml的 Project 部分、在上一节所述依赖之外显式添加如下依赖dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-google-cloud-dataflow-java/artifactId version${beam.version}/version scoperuntime/scope /dependency然后在 Maven JAR 插件中配置mainClass名称plugin 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执行mvn package -Pdataflow-runner后运行ls target假设你的artifactId为beam-examples、版本为1.0.0应能看到如下输出beam-examples-bundled-1.0.0.jar在 Cloud Dataflow 上运行该自执行 JAR 的命令如下java -jar target/beam-examples-bundled-1.0.0.jar \ --runnerDataflowRunner \ --projectYOUR_GCP_PROJECT_ID \ --regionGCP_REGION \ --tempLocationgs://YOUR_GCS_BUCKET/temp/ \ --outputgs://YOUR_GCS_BUCKET/output命令中各参数含义--runnerDataflowRunner指定执行后端--project指定 GCP 项目 ID--region指定作业创建所在的 Compute Engine 区域--tempLocation指定临时文件存放的 GCS 路径必须以gs://开头--output是管道业务参数示例中用于指定结果输出路径。从源码看Java 侧对jobName、project等提交参数有严格的校验逻辑作业名会被强制转换为小写并匹配[a-z](https://link.gitcode.com/i/65f9446d2b0b3a7e56c77ff7d32368e9)?正则项目 ID 必须是合法格式纯数字会被拒绝因为那是项目编号而非项目 ID见 DataflowRunner.java。Cloud Dataflow Runner 的 Pipeline 选项执行管道时需要为 Cloud Dataflow RunnerJava 或 Python考虑以下常用 Pipeline 选项。下表中的字段名采用 Java 风格camelCasePython 侧对应选项名已在表中标注。字段描述默认值runner使用的管道 Runner允许在运行时决定执行后端。设置为dataflow或DataflowRunner即可在 Cloud Dataflow 服务上运行。project你的 Google Cloud 项目 ID。若未设置默认取当前环境中的默认项目由gcloud配置。region创建作业的 Google Compute Engine 区域。若未设置默认取当前环境中的默认区域由gcloud配置。streaming是否启用流式模式处理无界PCollection时必须设为true。falsetempLocationPythontemp_location临时文件路径。必须是gs://开头的合法 Cloud Storage URL。Java 侧为可选Python 侧为必填。Java 中若设置了该值它将作为gcpTempLocation的默认值。无默认值。gcpTempLocation仅 Java临时文件的 Cloud Storage 存储桶路径必须是gs://开头的合法 URL。若未设置默认取tempLocation的值前提是它是合法的 Cloud Storage URL若tempLocation不是合法 Cloud Storage URL则必须显式设置gcpTempLocation。stagingLocationPythonstaging_location可选。用于暂存你的二进制文件和临时文件的 Cloud Storage 存储桶路径必须是gs://开头的合法 URL。Java未设置时默认为gcpTempLocation下的一个 staging 子目录Python未设置时默认为temp_location下的一个 staging 子目录。save_main_session仅 Python保存主会话状态使得定义在__main__中的函数和类例如交互式会话中定义的可以被 pickle 并在 worker 端反序列化。如果你的所有函数/类都定义在正式模块而非__main__中且模块在 worker 上可导入则某些工作流不需要该选项。falsesdk_location仅 Python覆盖 Beam SDK 的默认下载位置。可以是一个 URL、Cloud Storage 路径或本地 SDK tarball 路径。作业提交时会从该位置下载或复制 SDK tarball。若设为字符串default则使用标准 SDK 位置若为空则不复制任何 SDK。default关于这些默认值源码中有更精确的佐证project 默认值Java 侧DataflowPipelineOptions.getProject()标注了Default.InstanceFactory(DefaultProjectFactory.class)与Validation.Required未显式指定时由工厂从gcloud环境解析默认项目见 DataflowPipelineOptions.java。region 默认值同样通过DefaultGcpRegionFactory工厂从环境读取默认区域见 DataflowPipelineOptions.java。stagingLocation 默认值由StagingLocationFactory生成——当用户未设置时日志会提示 No stagingLocation provided, falling back to gcpTempLocation并把gcpTempLocation下追加/staging目录作为默认值如果gcpTempLocation缺失或不是合法 GCS 路径会抛出带明确提示的IllegalArgumentException见 DataflowPipelineOptions.java。streaming 相关行为在DataflowRunner.run()中若管道包含无界数据源/汇Runner 会自动把streaming置为true并追加流式引擎相关实验选项见 DataflowRunner.java。Java 侧更多选项可以参考DataflowPipelineOptions接口及其子接口的参考文档Python 侧对应PipelineOptions文档。接口声明中还包含大量未在上表列出的进阶选项例如isUpdate用同名作业替换现有管道并保留状态、getCreateFromSnapshot从快照创建作业、getTemplateLocation生成模板文件、getServiceAccount以指定服务账号运行作业、getFlexRSGoalFlexible Resource Scheduling 目标含SPEED_OPTIMIZED/COST_OPTIMIZED等完整定义见 DataflowPipelineOptions.java。附加信息与注意事项监控你的作业管道执行期间可以通过 Dataflow Monitoring Interface监控界面或 Dataflow Command-line Interface命令行界面监控作业进度、查看执行细节并接收管道结果更新。Java 侧提交成功后Runner 会在日志中打印监控控制台的访问链接以及 gcloud 取消命令例如To access the Dataflow monitoring console, please navigate to monitoring URL To cancel the job using the gcloud tool, run: gcloud dataflow jobs cancel jobId见 DataflowRunner.java。阻塞式执行如需阻塞直到作业完成可以对pipeline.run()返回的PipelineResult调用 Java 的waitToFinishPython 为wait_until_finish。Cloud Dataflow Runner 在等待期间会持续打印作业状态更新与控制台消息。需要注意当结果仍连接到活跃作业时在命令行按CtrlC并不会取消作业。要取消作业请使用 Dataflow Monitoring Interface 或 Dataflow Command-line Interfacegcloud dataflow jobs cancel命令。流式执行如果管道使用了无界数据源或数据汇必须将streaming选项设置为true。使用流式执行时需牢记以下三点流式管道不会自行终止除非用户显式取消。你可以从 Dataflow Monitoring Interface 取消流式作业或使用 Dataflow Command-line Interface 的gcloud dataflow jobs cancel命令。流式作业默认使用n1-standard-2或更高的 Compute Engine 机型。不要覆盖此设置因为n1-standard-2是运行流式作业所需的最低机型。流式执行的计费方式与批处理不同定价存在差异请参阅 Cloud Dataflow 官方定价文档。从源码实现看流式管道还会触发额外的自动化处理DataflowRunner.run()在检测到流式管道时不仅会开启streaming还会自动追加streaming_engine_state_tag_encoding_v2_supported等实验项从 2.75.0 版本起还会追加enable_streaming_engine_state_tag_encoding_v2并在统一 workeruseUnifiedWorker模式下自动启用 Streaming Engine见 DataflowRunner.java。小结Cloud Dataflow Runner 把 Beam 的批流统一编程模型与 Google Cloud 的托管服务结合起来你只需完成 GCP 项目、API、认证与存储桶的初始化在 Java 中声明 Runner 依赖必要时打包自执行 JAR然后通过一组 Pipeline 选项runner、project、region、streaming、tempLocation、stagingLocation等即可提交作业。对于需要长时间运行、自动扩缩容和动态负载均衡的大规模作业它是开箱即用的托管选择。本文涉及的 Runner 核心实现与选项定义可在本仓库的 DataflowRunner.java 与 DataflowPipelineOptions.javaJava以及 pipeline_options.pyPython中继续深入阅读。赞分享大数据批处理流处理数据工程【免费下载链接】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 管道DataflowRunner 完整实战指南在 Google Cloud Dataflow 上运行 Apache Beam 管道DataflowRunner 完整实战指南 Apache Beam 提供统大数据批处理流处理数据工程Apache Hop 可视化管道上云实战使用 Google Cloud Dataflow 运行 Beam 管道完整指南Apache Hop 可视化管道上云实战使用 Google Cloud Dataflow 运行 Beam 管道完整指南 本文基于 Apache Beam 官方大数据批处理流处理数据工程Apache Beam 使用 Cloud Dataflow Runner 执行批流管道从项目配置、依赖打包到任务监控的完整实战指南Apache Beam 使用 Cloud Dataflow Runner 执行批流管道从项目配置、依赖打包到任务监控的完整实战指南 Apache Beam 通大数据批处理流处理数据工程上一篇Flet 路由高亮实战使用 is_route_active() 实现导航菜单激活状态下一篇RedwoodRecord 实战指南基于 Prisma 的 Redwood 原生 ORM 全解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
📌 标签:
工业官网
设计趋势
AI 建站
SEO
获取完整报告 →
RELATED ARTICLES
推荐阅读
2026/10/10 2:43:50
华三设备模拟MV专线-客户两端设备运行BGP
2026/10/10 2:38:50
markdown-it 基准测试样本解析:block-bq-flat.md 与扁平引用块的解析与压测原理
2026/10/10 2:38:50
深入理解 containers/storage:Go 语言实现的 Layer、Image、Container 容器存储层管理库
2026/10/10 3:38:54
Linux 7z 命令实战:高压缩比参数调优与避坑指南
2026/10/10 3:38:54
OpenClaw 配置教程:从零到跑通模型、技能与数据库
2026/10/10 3:38:54
RabbitMQ实战:从消息队列原理到分布式系统解耦与故障排查
2026/10/10 3:38:54
Java毕设实战:彝族文化宣传网站完整设计与实现方案
2026/10/10 3:38:54
Flutter三方库鸿蒙化适配实战:以波斯语本地化库persian为例
2026/10/10 3:33:54
计算机单片机毕设实战-基于单片机的室内多环境因子阈值配置与超标联动告警系统设计 基于单片机的密闭空间烟雾与甲醛监测及可控排风装置设计(030115)
2026/10/10 0:03:38
工业软件标准化路线图:国产替代的落地施工图
2026/10/10 0:03:38
VCMI安卓版实操指南:原生运行英雄无敌3的3步技术落地
2026/10/10 0:03:38
稀疏多通道盲反褶积的MATLAB算法实现与参数调优
2026/10/10 3:42:06
Jev+Agent接管浏览器:browser-use实战与jev-ultrafast性能优化
2026/10/10 3:42:01
多智能体集群实战:DeepAgents编排、MCP与A2A协议及Skills体系
2026/10/10 3:41:58
hindsight:面向LLM应用的事后可观测性工程实践
2026/10/10 3:41:56
我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频
2026/10/10 3:41:54
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证
2026/10/9 11:36:17
2026 大模型集体涨价:用 Python 做企业 Token 成本测算与选型避坑(附配置)