- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文基于 Apache Beam 官方仓库博客 Running Apache Hop visual pipelines with Google Cloud Dataflow 整理。Apache Hop 是一款面向 Apache Beam 的可视化管道开发环境,它让开发者无需编写代码即可设计、预览、调试和部署 Beam 管道,并可在 Spark、Flink、Google Cloud Dataflow 等任意 Beam Runner 上运行。读完本文,你将掌握从安装 Hop、本地预览管道、使用 Beam Direct Runner 本地运行,到配置 Google Cloud 环境、生成 fat jar、最终在 Dataflow 云端规模化运行同一套可视化管道的完整实战流程。
一、Apache Hop:Beam 的可视化管道开发环境
Apache Hop 是 Apache 软件基金会旗下的可视化数据处理与编排平台,它把 Apache Beam 的"一次设计,随处运行"(design once, run anywhere)理念延伸到可视化管道开发与生命周期管理上。根据仓库中的案例研究 Visual Apache Beam Pipeline Design and Orchestration with Apache Hop,Hop 使用元数据(metadata)和内核(kernel)描述数据应如何处理,再由 Apache Beam 负责在底层执行引擎上运行;其 GUI 允许开发者以拖拽方式可视化地构建管道,所有管道元素的设置只需在 Hop 可视化编辑器中配置一次,管道即自动以 JSON、CSV 等元数据格式描述,无需掌握特定编程语言即可创建管道。
Hop 的核心工作方式可以概括为:
- Transform(变换):管道中的数据处理节点,例如读入 CSV、过滤、选择列等;
- Hop 连接线:将各个 Transform 链接起来,形成完整的管道数据流;
- 管道运行配置(Pipeline Run Configuration):指定管道在哪个 Beam Runner(Direct、Spark、Flink、Dataflow 等)上执行。
本教程以一个随 Hop 客户端分发的示例管道input-process-output为主线,演示完整的"本地预览 → 本地运行 → 云端运行"三部曲。由于 Hop 的 Beam 管道本质上是 Beam 管道,因此最终部署到 Dataflow 时,只差"为使用 Dataflow 设置必要的参数"这一步。
二、安装与启动 Apache Hop
Apache Hop 既可以作为本地应用程序运行,也可以使用 Docker 容器运行 Hop Web 版本。本文的指令针对本地应用版本,因为容器环境下 Cloud Dataflow 的认证方式会有所不同,其余所有指令与界面操作在两种形态下完全一致(Hop 的 UI 相同)。
2.1 下载与解压
下载并安装 Apache Hop 后即可开始使用。本文使用apache-hop-client包中的二进制文件,版本为1.2.0(发布于 2022 年 3 月 7 日)。解压后的 Zip 包内包含一个config目录,其中存放着:
- 若干示例项目(sample projects);
- 针对 Dataflow 及其他 Runner 的管道运行配置。
本文使用的示例管道位于config/projects/samples/beam/pipelines/input-process-output.hpl。
2.2 启动 Hop
在解压客户端所在目录下运行启动脚本:
./hop/hop-gui.shWindows 用户则运行:
./hop/hop-gui.bat2.3 打开示例管道
- 启动后,将左上角
projects框中的项目从default切换为samples; - 点击打开按钮(Open);
- 在
beam/pipelines子目录中选择管道input-process-output.hpl。
此时主窗口会显示类似下图的管道图:Customers输入节点读取 CSV 客户数据,Only CA节点过滤出stateCode列等于CA的记录,Limit fields, re-order节点仅保留部分列并重排列顺序,最后将结果写入 Google Cloud Storage。简言之,这条管道完成了"过滤 → 选列 → 输出到 GCS"的典型 ETL 流程。
三、本地预览每个变换的输出
在上云之前,先在本地验证管道逻辑是明智的做法。Hop 支持对每个 Transform 预览输出:
- 点击输入节点
Customers; - 在弹出的对话框中选择Preview Output(预览输出);
- 选择Quick launch(快速启动),即可看到部分输入数据;
- 查看完毕后点击Stop停止预览。
同样地,在Only CA变换之后重复上述操作,可以验证所有行的stateCode列都等于CA,说明过滤逻辑生效。
下一个变换Limit fields, re-order会从输入数据中仅选择部分列并重新排序。点击该变换并选择Preview Output,再点击Quick Launch,可以看到输出中id列被排到了第一位,且只保留了输入列的一个子集——这正是管道写完完整输出后数据的样子。这一步的预览结果,就是你期望 Dataflow 云端作业最终产出的结果。
四、使用 Beam Direct Runner 本地运行
预览验证通过后,可以在本地完整运行一次管道,确认端到端流程可用。
4.1 认识管道运行配置
运行管道需要指定一个 Runner 配置,这通过 Apache Hop 的Metadata 工具完成。在samples项目中,已经预置了多个运行配置:
local配置:用于在 Hop 内运行管道(例如前面预览各步骤输出时使用的配置);Direct配置:使用 Apache Beam 的 Direct Runner。
管道运行配置(Pipeline Run Configurations)包含两个页签:main和variables。
**Main 页签(Direct Runner)**主要包含 worker 数量等设置。你可以根据本机 CPU 数量调整 worker 数,甚至可以限制为 1,让管道不消耗过多资源。值得注意的是,Direct Runner 输出文件的数量与运行配置中设置的 worker 数量相关。
Variables 页签则保存管道自身的配置参数(而非 Runner 的)。对于本示例管道,只使用了DATA_INPUT和DATA_OUTPUT两个变量;另一个STATE_INPUT变量用于不同的示例。你可以打开管道输入、输出节点的 Beam 变换,看到这些变量是如何在节点中引用的——这正是 Hop 管道"参数化"能力的体现:通过变量把数据路径与管道逻辑解耦。
4.2 运行管道
由于DATA_INPUT、DATA_OUTPUT变量已正确指向 samples 项目文件夹中的数据位置,可以直接尝试用 Beam Direct Runner 运行:
- 回到管道视图(Metadata 工具上方的箭头按钮);
- 点击工具栏中小的"播放"按钮(run);
- 选择
Direct管道运行配置; - 点击Launch按钮。
如何判断作业是否完成?查看主窗口底部的日志即可。日志中会显示管道使用 "Beam Pipeline Engine" 执行并最终完成的信息。
运行完成后,前往DATA_OUTPUT变量指定的位置——本例为config/projects/samples/beam/output——可以看到若干输出文件。输出文件的数量取决于运行配置中设置的 worker 数。
补充:为什么先跑 Direct Runner?根据仓库中的 Using the Direct Runner 文档,Direct Runner 在本机执行管道,旨在尽可能严格地验证管道是否符合 Beam 模型,例如强制元素不可变性、可编码性、元素在任意点的任意处理顺序、用户函数(
DoFn、CombineFn)的可序列化性等。它面向正确性而非性能,需要把所有用户数据放入内存,因此不适合生产管道,却是开发与单元测试阶段验证管道健壮性的首选。在 Direct Runner 上验证通过,意味着管道语义正确,可以放心迁移到远程 Runner。
本地运行成功,接下来就是上云时刻。
五、为云端运行准备 Google Cloud 环境
要跟随本文在云端运行,你需要具备:一个 Google Cloud Platform 项目、足够的权限创建(或使用已有)Google Cloud Storage 存储桶、以及运行 Dataflow 作业的权限。
5.1 检查并配置 Google Cloud SDK
首先安装 Google Cloud SDK(gcloud),并配置其使用你的账号与项目。
双重检查配置是否正确:
gcloud config list确认输出的 account 和 project 无误后,再做一次认证:
gcloud auth login该命令会在浏览器中打开标签页完成认证。认证完成后,即可通过 SDK 与你的项目交互。
5.2 创建 GCS 存储桶并上传示例数据
本文示例使用 GCP 的europe-west1区域。创建一个区域级存储桶(bucket 名称务必自行更换,避免与他人冲突):
gcloud storage buckets create gs://YOUR-BUCKET-NAME --location=europe-west1 --default-storage-class=standard接下来把示例数据上传到 GCS 存储桶,以测试管道在 Dataflow 中的运行效果。进入存放所有 Hop 文件的目录(即hop-gui.sh所在的同一目录),复制数据到 GCS:
gcloud storage cp config/projects/samples/beam/input/customers-noheader-1k.txt gs://YOUR-BUCKET-NAME/data/注意路径末尾的斜杠/:它表示要创建一个名为data的目录,并把所有内容放入其中。
验证数据上传是否正确:
gcloud storage ls gs://YOUR-BUCKET-NAME/data/该位置应能看到文件customers-noheader-1k.txt。
5.3 启用 Dataflow 并准备凭据
继续之前,请确保项目中已启用 Dataflow,并准备好可与 Hop 配合使用的服务账号。可以参考 Dataflow 官方文档的 "Before you begin" 章节完成 API 启用。
Hop 官方文档推荐的认证方式是:创建一个服务账号、为其导出 key、并设置GOOGLE_APPLICATION_CREDENTIALS环境变量。但导出服务账号 key 存在潜在安全风险,因此本文改用更安全的另一种方式——借助 Google Cloud SDK:
gcloud auth application-default login该命令会在浏览器中打开标签页要求确认认证。确认后,系统上任何需要访问 Google Cloud Platform 的应用程序都会使用这份凭据。
5.4 创建 Dataflow 服务账号并授权
还需要为 Dataflow 作业创建一个具有特定权限的服务账号:
gcloud iam service-accounts create dataflow-hop-sa然后为这个服务账号授予 Dataflow 相关权限(注意把YOUR-PROJECT-ID替换为你的项目 ID):
gcloud projects add-iam-policy-binding YOUR-PROJECT-ID \ --member="serviceAccount:dataflow-hop-sa@YOUR-PROJECT-ID.iam.gserviceaccount.com" \ --role="roles/dataflow.worker"还需要授予 Google Cloud Storage 的额外权限:
gcloud projects add-iam-policy-binding YOUR-PROJECT-ID \ --member="serviceAccount:dataflow-hop-sa@YOUR-PROJECT-ID.iam.gserviceaccount.com" \ --role="roles/storage.admin"最后,授权你的用户去模拟(impersonate)该服务账号:
- 在 Google Cloud Console 中打开项目的 Service Accounts 页面;
- 点击刚创建的服务账号;
- 点击Permissions页签,再点击Grant Access按钮;
- 给你的用户授予Service Account User角色。
至此,你的用户就具备使用该服务账号运行 Dataflow 的全部前置条件了。
六、更新管道运行配置并启动 Dataflow 作业
6.1 生成 fat jar
在 Dataflow 上运行管道之前,需要为管道代码生成 JAR 包:
- 在菜单栏的Tools菜单中,选择Generate a Hop fat jar;
- 在对话框中点击 OK;
- 选择 JAR 的位置与文件名,点击Save。
生成该文件需要几分钟时间。
6.2 进入 Dataflow 管道运行配置
进入管道编辑器,点击播放按钮,选择DataFlow作为管道运行配置,再点击右侧的播放按钮,即可打开 Dataflow 管道运行配置界面,在这里可以修改输入变量及其他 Dataflow 设置。
6.3 修改 Variables 页签
点击Variables页签,只修改DATA_INPUT和DATA_OUTPUT两个变量——注意,这里不仅路径要指向你自己的 GCS 存储桶,文件名也需要一并修改。默认的DATA_INPUT、DATA_OUTPUT是指向示例项目作者 GCS 位置的,必须改为你自己的存储桶地址。
6.4 修改 Main 页签
接下来进入Main页签,需要更新的选项包括:
| 配置项 | 说明 |
|---|---|
| Project id | 你的 GCP 项目 ID(即gcloud config list中显示的 project 值) |
| Service account | 刚创建的服务账号地址(可在 Console 的 Service Accounts 页面找到) |
| Staging location | 暂存位置,使用刚创建的存储桶,只改桶地址,保留配置中已有的 "binaries" 子路径 |
| Region | 本例使用europe-west1 |
| Temp location | 临时文件位置,同样使用刚创建的存储桶,保留 "tmp" 子路径 |
| Fat jar file location | 点击右侧Browse按钮,定位之前生成的 JAR 文件 |
其中 staging 与 temp 位置都指向同一个存储桶,仅需把路径中的桶地址换成你自己的,其余子路径("binaries"、"tmp")保持不变。
此外,根据你的网络配置,可能需要勾选"Use Public IPs?"复选框;或者不勾选,但需要在你的项目europe-west1区域的子网中启用 Private Google Access(即 VPC 内的私有 Google 访问)。本文为简单起见选择勾选 Use Public IPs。
你当然也可以根据项目具体要求调整任何其他选项。
6.5 理解这些参数背后的 Dataflow Runner 语义
Hop 的 Main 页签参数与 Beam 的 Dataflow Runner 管道选项一一对应。根据仓库中的 Using the Google Cloud Dataflow Runner 文档,关键选项及其默认行为如下:
runner:指定运行 Runner,设为dataflow或DataflowRunner即运行于 Cloud Dataflow 托管服务;project:Google Cloud 项目 ID,未设置时默认取当前环境(gcloud配置)中的默认项目;region:创建作业的 Compute Engine 区域,未设置时默认取当前环境的默认区域;streaming:是否启用流式模式(处理无界PCollection时需设为true),默认false;本文的批处理管道无需开启;tempLocation/gcpTempLocation:临时文件路径,必须是gs://开头的合法 Cloud Storage URL;stagingLocation:存放二进制与临时文件的暂存路径,必须是gs://开头的 URL,未设置时默认落在gcpTempLocation内的 staging 目录。
从这些语义可以理解:Dataflow Runner 会把可执行代码与依赖上传到 GCS 存储桶,并创建 Cloud Dataflow 作业,在 Google Cloud 的托管资源上执行管道。这正是 Hop 配置中 Staging/Temp location 必须指向gs://路径的原因。
6.6 启动作业
配置完成后,点击Ok按钮,再点击Launch触发管道。在日志窗口中,你应该能看到类似"作业已提交"的一行输出,说明 Hop 已成功把管道提交给了 Dataflow。
七、在 Dataflow 中检查作业
一切顺利的话,现在可以在 Google Cloud Console 的 Dataflow 作业页面(Jobs 页面)看到一个正在运行的作业。
- 查看作业列表:作业页会列出
input-process-output批处理作业及其运行状态; - 排查失败:如果作业失败,打开失败作业页面,查看底部的Logs,点击错误图标即可定位失败原因——绝大多数情况是运行配置中某项参数设置错误;
- 查看执行图:管道开始运行后,作业页面会展示管道的执行图,对应 Hop 中的各个变换节点(Customers、Only CA、Limit fields 等)逐个进入 Running 状态。
7.1 校验输出
作业结束后,输出位置应有文件生成,用gcloud storage查看:
gcloud storage ls gs://YOUR-BUCKET-NAME/output示例输出(实际文件名与数量每次运行可能不同):
gs://YOUR-BUCKET-NAME/output/input-process-output-00000-of-00003.csv gs://YOUR-BUCKET-NAME/output/input-process-output-00001-of-00003.csv gs://YOUR-BUCKET-NAME/output/input-process-output-00002-of-00003.csv读取输出文件的前几行:
gcloud storage cat "gs://YOUR-BUCKET-NAME/output/*csv" | head示例输出如下:
12,wha-firstname,vnaov-name,egm-city,CALIFORNIA 25,ayl-firstname,bwkoe-name,rtw-city,CALIFORNIA 26,zio-firstname,rezku-name,nvt-city,CALIFORNIA 44,rgh-firstname,wzkjq-name,hkm-city,CALIFORNIA 135,ttv-firstname,eqley-name,trs-city,CALIFORNIA 177,ahc-firstname,nltvw-name,uxf-city,CALIFORNIA 181,kxv-firstname,bxerk-name,sek-city,CALIFORNIA 272,wpy-firstname,qxjcn-name,rew-city,CALIFORNIA 304,skq-firstname,cqapx-name,akw-city,CALIFORNIA 308,sfu-firstname,ibfdt-name,kqf-city,CALIFORNIA可以看到:所有行的州(state)都是 CALIFORNIA(即CA过滤生效);输出只包含选定的列;用户 id 是第一列。你实际得到的输出大概率与此不同,因为每次运行的数据处理顺序不尽相同。
7.2 规模扩展
本次作业只使用了小样本数据,但同一个作业完全可以替换为任意大规模输入 CSV——Dataflow 会自动并行、分块处理数据,这正是把 Hop 可视化管道迁移到云端托管服务获得的核心价值。
八、总结与延伸
Apache Hop 是 Beam 管道的可视化开发环境,支持本地运行管道、检查数据、调试、单元测试等丰富能力。一旦管道在本地验证通过,只需设置使用 Dataflow 所需的必要参数,就可以把同一套可视化管道直接部署到云端,实现从"本地可用"到"云端规模化"的平滑过渡。
进一步学习路径(均位于本仓库):
- 如果不想在本机安装任何软件,可参考姊妹教程 Apache Hop web version with Cloud Dataflow:它演示了如何在 Google Cloud 虚拟机中运行 Hop Web 容器、通过 Identity Aware Proxy 安全访问,并用浏览器直接创建管道、启动 Dataflow 作业,其架构图与部署命令均可直接复现;
- 想了解 Hop 项目背景与"可视化 Beam 管道"的设计理念,可阅读案例研究 Visual Apache Beam Pipeline Design and Orchestration with Apache Hop;
- 想深入掌握 Dataflow Runner 的完整管道选项(project、region、streaming、tempLocation、gcpTempLocation、stagingLocation 等),可查阅 Using the Google Cloud Dataflow Runner;
- 想理解 Direct Runner 为何适合本地验证(模型语义检查、内存限制、并行度设置等),可查阅 Using the Direct Runner。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
在 Google Cloud Dataflow 上运行 Apache Beam 管道:DataflowRunner 完整实战指南
在 Google Cloud Dataflow 上运行 Apache Beam 管道:DataflowRunner 完整实战指南 Apache Beam 提供统
大数据批处理流处理数据工程Google Cloud Dataflow 实战指南:基于 Apache Beam 的流批一体数据管道服务
Google Cloud Dataflow 实战指南:基于 Apache Beam 的流批一体数据管道服务 Google Cloud Dataflow 是 Go
云原生CI/CD运维在 Google Cloud Dataflow 上运行 Apache Beam 模板:python-docs-samples run_template 实战指南
在 Google Cloud Dataflow 上运行 Apache Beam 模板:python docs samples run_template 实战指南
示例工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考