ARTICLE DETAIL

资讯详情

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

Apache Beam 代码修改实战指南:从 Gradle 构建体系到本地测试与产物发布

Apache Beam 代码修改实战指南:从 Gradle 构建体系到本地测试与产物发布 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本文基于 Apache Beam 仓库的官方贡献文档contributor-docs/code-change-guide.md展开系统讲解 Beam monorepo 的目录结构、Gradle 构建体系、Java/Python SDK 的本地测试流程以及如何用修改后的 Beam 代码运行自己的 Pipeline。读完后你可以独立完成在本地修改 Beam 源码、运行单元/集成测试、将改动发布到 Maven Local 或构建 SDK 容器镜像并在 Dataflow 等 Runner 上验证自己的修改。仓库结构一个覆盖全栈的 monorepoApache Beam 的 GitHub 仓库在很大程度上是一个 monorepo它包含 Beam 项目的一切包括各语言 SDK、测试基础设施、数据看板dashboards、Beam 官网和 Beam Playground。对代码修改工作而言理解以下关键目录是第一步。代码路径速览Java 代码路径Java 代码主要集中在sdks/java和runners两个目录sdks/java—— Java SDK 主体sdks/java/core—— Java 核心Pipeline、PCollection、Coder 等基础组件sdks/java/harness—— SDK harness即 SDK 容器的入口点runners—— Java Runner 实现包括runners/direct-java—— Direct Runner本地 Runner测试首选runners/flink—— Flink Runner当前仓库中按 Flink 大版本划分如runners/flink/2.0、runners/flink/2.2等runners/google-cloud-dataflow-java—— Dataflow Runner作业提交、翻译等runners/google-cloud-dataflow-java/worker—— 供不使用 Runner v2 的 Dataflow 作业使用的 Worker非 Java SDK 路径其他语言的 SDK 位于sdks/LANG目录重点子目录说明sdks/python—— Python 包的 setup 文件与触发 test-suites 的脚本sdks/python/apache_beam—— Beam Python 包本体sdks/python/apache_beam/runners/worker—— SDK worker harness 入口与 state samplersdks/python/apache_beam/io—— I/O 连接器sdks/python/apache_beam/transforms—— 大部分核心组件sdks/python/apache_beam/ml—— Beam ML 代码sdks/python/apache_beam/runners—— Runner 实现与封装sdks/go—— Go SDKGo SDK 的测试与变更指引见仓库内 sdks/go/README.md.github/workflows—— GitHub Actions 工作流即 PR 阶段运行的测试。绝大多数工作流只执行一条 Gradle 命令。开发本地测试时可以查看对应工作流中运行的命令以确认本地该执行什么。Gradle 快速上手Beam 仓库是一个单一的 Gradle 工程包含 Java、Python、Go 和 website 在内的所有组件。开始开发前建议先熟悉 Gradle 多项目构建的基本概念project / task / pluginproject包含build.gradle文件的目录task在build.gradle中定义的动作plugin预定义的任务与任务层级在项目build.gradle中生效Java 项目或子项目的常见任务任务作用compileJava编译 Java 源码compileTestJava编译 Java 测试源码test运行单元测试integrationTest运行集成测试执行任务有两种等价写法./gradlew -p sdks/java/core compileJava ./gradlew :sdks:java:harness:testBeam 专属的 Gradle 插件BeamModulePlugin对 Apache Beam 而言几乎所有工程配置都由一个插件统一管理BeamModulePlugin。从源码结构看该插件BeamModulePlugin.groovy#L65向项目注入了一组被称为 natures 的配置方法覆盖 Java、Go、Docker、Grpc、Avro 等项目类型职责包括管理 Java 依赖为 Java、Python、Go、Proto、Docker、Grpc、Avro 等项目配置标准插件Java 项目调用applyJavaNature源码见 BeamModulePlugin.groovy#L1146Python 项目调用applyPythonNature源码见 BeamModulePlugin.groovy#L3148为每种项目类型定义公共自定义任务例如test运行 Java 单元测试、spotlessApply格式化 Java 代码因此仓库中每个 Java 项目的build.gradle开头都是这样的模式apply plugin: org.apache.beam.module applyJavaNature( ... )理解了这一点后你在任意子目录看到的行为差异如额外的docker、integrationTest任务都能在该插件中追溯到来源。环境准备设置本地开发环境前请先阅读仓库根目录的 CONTRIBUTING.md 贡献指南。如果计划使用 Dataflow还需要配置gcloud凭据。根据涉及的语言PATH中需要配置以下环境Java 环境使用受支持的 Java 版本推荐 Java 11。由于 Beam 是基于 JVM 的 Gradle 工程所有开发都需要该环境。推荐用 sdkman 管理 Java 版本。Python 环境任意受支持的 Python 版本仅 Python SDK 开发需要。推荐pyenv 虚拟环境venv管理。Go 环境最新 Go 版本Go SDK 开发需要由于所有 SDK 的容器 entrypoint 脚本均用 Go 编写修改任意 SDK 容器也需要 Go。Docker 环境SDK 容器变更、部分跨语言功能如果本地运行 SDK 容器镜像Beam 2.53.0 及以后版本不再必需、以及 Flink/Spark 等 portable runner 的 job server 场景都需要。判断所需环境的经验法则修改sdks/java/io/google-cloud-platform的代码 → 只需 Java 环境修改sdks/java/harness的代码 → 需要 Java Go DockerDocker 用于编译构建 Java SDK harness 容器镜像修改sdks/python/apache_beam的代码 → 需要 Python 环境。Java 开发指南IntelliJ 设置IDE 并非必需命令行同样可以完成修改与测试但如果使用 IntelliJ注意以下步骤用 IntelliJ 打开仓库根目录/beam重要打开根目录而不是sdks/java。等待索引完成可能需要几分钟。由于 Gradle 是自包含的构建工具满足前置条件后环境设置即完成。验证加载是否成功找到examples/java/build.gradle文件在wordCount任务旁会出现Run按钮点击后 wordCount 示例会编译并运行。从源码看该任务的定义位于 examples/java/build.gradle#L161-L167它是一个JavaExec任务主类为org.apache.beam.examples.WordCount参数固定为--output/tmp/output.txt并使用当前系统属性。因此该任务同时被用作贡献者 Java 环境自检。命令行环境验证用 Gradle 命令行运行测试时推荐先执行下面这条命令。它会编译 Apache Beam SDK、WordCount Pipeline 及一个数据处理 Hello-world 程序并在 Direct Runner 上运行该 Pipeline$ cd beam $ ./gradlew :examples:java:wordCount成功时 Gradle 构建日志会出现类似输出... BUILD SUCCESSFUL in 2m 32s 96 actionable tasks: 9 executed, 87 up-to-date 3:41:06 PM: Execution finished wordCount.同时输出文件/tmp/output.txt-*中出现词频统计结果例如$ head /tmp/output.txt* /tmp/output.txt-00000-of-00003 should: 38 bites: 1 depraved: 1 gauntlet: 1 battle: 6 ...如果该命令能走通说明 Java Gradle 环境已就绪。运行单元测试假设你在 Java SDK 的某个模块例如sdks/java/io/jdbc做了代码修改接下来本地跑单元测试测试统一存放在各项目的src/test/java目录单元测试文件名形如.../**Test.java集成测试形如.../**IT.java。运行整个项目的单元测试./gradlew :sdks:java:harness:testJUnit HTML 报告位于invoked_project/build/reports/tests/test/index.html。运行特定测试支持全限定类名、通配类名、具体方法三种粒度./gradlew :sdks:java:harness:test --tests org.apache.beam.fn.harness.CachesTest ./gradlew :sdks:java:harness:test --tests *CachesTest ./gradlew :sdks:java:harness:test --tests *CachesTest.testClearableCache在 IntelliJ 中点击测试类或方法旁的绿色三角即可运行/调试设置断点即可 debug。特例sdks:java:core的单元测试不走上述test任务。由于 Java core 不依赖任何 Runner凡是需要真正运行 Pipeline 的单元测试必须借助 Direct Runner因此要使用:runners:direct-java:needsRunnerTest来触发 core 的相关测试。运行集成测试集成测试文件名形如.../**IT.java基于 TestPipeline位于sdks/java/core的org.apache.beam.sdk.testing包编写通过TestPipelineOptions设置选项。集成测试与标准 Pipeline 的区别默认在run时阻塞等待默认使用TestDataflowRunner默认超时 15 分钟Pipeline 选项通过系统属性beamTestPipelineOptions注入。通过-DbeamTestPipelineOptions[...]可以显式配置测试使用的 Pipeline 选项例如-DbeamTestPipelineOptions[--runnerTestDataflowRunner,--projectmygcpproject,--regionus-central1,--stagingLocationgs://mygcsbucket/path]build.gradle 中预配置的 beamTestPipelineOptions某些项目在build.gradle中显式组装了该属性。以 sdks/java/io/google-cloud-platform/build.gradle 为例其integrationTest任务从项目属性gcpProject、gcpTempRoot等取值未赋值时默认回退到apache-beam-testingGCP 项目temp root 默认为gs://temp-storage-for-end-to-end-tests并用include **/*IT.class只匹配集成测试类同时通过outputs.upToDateWhen { false }禁用 Gradle 缓存——因为集成测试访问真实服务始终视为 out of date。要在自己的项目下运行./gradlew :sdks:java:io:google-cloud-platform:integrationTest -PgcpProjectmygcpproject -PgcpTempRootgs://mygcsbucket/path另一些项目如sdks/java/io/jdbc、sdks/java/io/kafka不在build.gradle中组装beamTestPipelineOptions直接用-DbeamTestPipelineOptions[...]显式传入即可。编写集成测试在集成测试中设置TestPipeline对象的标准写法Rule public TestPipeline pipelineWrite TestPipeline.create(); Test public void testSomething() { pipeline.apply(...); pipeline.run().waitUntilFinish(); }运行测试的任务需要指定 Runner常见任务映射Google Cloud I/O 集成测试跑 Direct Runner:sdks:java:io:google-cloud-platform:integrationTest标准 Dataflow Runner:runners:google-cloud-dataflow-java:googleCloudPlatformLegacyWorkerIntegrationTestDataflow Runner v2:runners:google-cloud-dataflow-java:googleCloudPlatformRunnerV2IntegrationTest查看 GitHub Action 工作流实际执行的 Gradle 命令即可了解 CI 是如何本地复现这套流程的。一个典型调用示例./gradlew :runners:google-cloud-dataflow-java:examplesJavaRunnerV2IntegrationTest \ -PdisableSpotlessChecktrue -PdisableCheckStyletrue -PskipCheckerFramework \ -PgcpProjectyour_gcp_project -PgcpRegionus-central1 \ -PgcsTempRootgs://your_gcs_bucket/tmp用修改后的 Beam 代码运行 Pipeline将代码改动应用到自己的 Pipeline 时建议先从独立分支开始提交 PR 或在 dev 分支上验证修改 → 基于 Beam HEADmaster基于已发布版本2.xx.0打补丁 → 基于发布 tag如v2.55.0。随后在 Beam 仓库中编译包含修改的项目并发布到 Maven Local以修改sdks/java/io/kafka为例./gradlew -Ppublishing -p sdks/java/io/kafka publishToMavenLocal该命令默认把修改后的构件发布到 Maven Local 仓库~/.m2/repository用户 Pipeline 运行时即可拾取该变更。SNAPSHOT 版本的处理如果改动位于 master 或 PR 等开发分支产物版本为2.xx.0-SNAPSHOT而非发布 tag 版本Pipeline 项目需要额外配置仓库来源。Maven 项目步骤推荐以 WordCount 的maven-archetype模板搭建项目添加 snapshot 仓库repository idMaven-Snapshot/id namemaven snapshot repository/name urlhttps://repository.apache.org/content/groups/snapshots//url releases enabledfalse/enabled /releases /repository在pom.xml中将beam.version改为对应的快照版本properties beam.version2.XX.0-SNAPSHOT/beam.version /propertiesGradle 项目步骤在build.gradle中加入仓库配置repositories { mavenCentral() mavenLocal() maven { url https://repository.apache.org/content/groups/snapshots } }将 Beam 依赖版本设为2.XX.0-SNAPSHOT。该配置会让构建系统从 Maven Snapshot 仓库下载 Beam 的 nightly 构建而不会下载你本地修改的构建。通常不需要在本地构建全部 Beam 产物确有需要时执行./gradlew -Ppublishing publishToMavenLocal标准 Dataflow Runner非 v2且 worker harness 有改动编译dataflowWorkerJar./gradlew :runners:google-cloud-dataflow-java:worker:shadowJarjar 位于构建输出目录beam root/runners/google-cloud-dataflow-java/worker/build/libs。在 Pipeline 提交参数中加入--dataflowWorkerJarfull-path-to/beam-runners-google-cloud-dataflow-java-legacy-worker-2.XX.0-SNAPSHOT.jarDataflow Runner v2 且sdks/java/harness或其依赖如sdks/java/core有改动构建 SDK harness 容器./gradlew :sdks:java:harness:publishToMavenLocal -Ppublishing在项目依赖中加入org.apache.beam:beam-sdks-java-harness:2.XX.0-SNAPSHOTGradle 示例implementation(org.apache.beam:beam-sdks-java-harness:2.XX.0-SNAPSHOT)提交 Pipeline 时附加实验参数--experimentsuse_runner_v2,use_staged_dataflow_worker_jarSDK 容器镜像的修改如果你修改了sdks/java/container下的代码必须构建自定义 SDK 容器才能生效./gradlew :sdks:java:container:java11:docker # 或 java17、java21 等 # 将版本 tag 改为实际值 docker tag apache/beam_java8_sdk:2.68.0.dev \ us-docker.pkg.dev/apache-beam-testing/beam-temp/beam_java11_sdk:2.68.0-custom # 改成你的 artifact registry docker push us-docker.pkg.dev/apache-beam-testing/beam-temp/beam_java11_sdk:2.68.0-custom然后用以下选项运行 Pipeline--experimentsuse_runner_v2 \ --sdkContainerImageus.gcr.io/apache-beam-testing/beam_java11_sdk:2.49.0-customSnapshot 版本容器开发中的 SDK snapshot 版本默认使用发布到 apache-beam-testing 项目容器注册表中的镜像https://gcr.io/apache-beam-testing/beam-sdk/...例如最新的 Java 21 snapshot 容器即位于该注册表下。当版本进入发布候选阶段时会再发布最后一个 SNAPSHOT 版本该版本改用发布在 DockerHub 上的最终容器。需要注意发布流程期间 SNAPSHOT 容器可能出现短暂不可用。为避免影响建议切换到可用的最新 SNAPSHOT 版本或使用自定义容器重要工作负载尽量不要长期依赖 snapshot 版本。某些 Runner 会覆盖此 snapshot 行为例如 Dataflow Runner 会把所有 SNAPSHOT 容器映射到单一注册表gcr.io/cloud-dataflow/切换到最终容器时同样会经历同样的短暂不可用窗口。Python 开发指南Beam Python SDK 以单个 wheel 分发比 Java SDK 简单直接得多。控制台环境配置以下以工作目录sdks/python为例推荐用pyenv安装 Python 解释器安装前置依赖curl https://pyenv.run | bashpyenv install 3.X受支持的 Python 版本参考 gradle.properties 中的python_version项目属性创建并激活虚拟环境pyenv virtualenv 3.X ENV_NAMEpyenv activate ENV_NAME以 editable 模式安装apache_beam包cd sdks/python pip install -e .[gcp,test]若涉及 SDK 容器镜像的开发还需安装 Docker Desktop 与 Go。如果打算提交 PR安装 pre-commit 钩子避免 Python 代码的 lint 失败# 启用 pre-commit (env) $ pip install pre-commit (env) $ pre-commit install # 禁用 pre-commit (env) $ pre-commit uninstall运行单元测试虽然测试也可以用 Gradle 命令触发但那会在每次运行前新建virtualenv并安装依赖耗时数分钟因此持久化虚拟环境更有价值。单元测试文件名形如**_test.py。运行整个文件的全部测试pytest -v apache_beam/io/textio_test.py运行某个测试类的全部测试pytest -v apache_beam/io/textio_test.py::TextSourceTest运行单个测试方法pytest -v apache_beam/io/textio_test.py::TextSourceTest::test_progress运行集成测试集成测试文件名形如**_it_test.py。在 Direct Runner 上运行集成测试python -m pytest -o log_cliTrue -o log_levelInfo \ apache_beam/ml/inference/pytorch_inference_it_test.py::PyTorchInference \ --test-pipeline-options--runnerTestDirectRunner如果准备提交 PR希望测试进入 PostCommit Python 的 test-suites需要在 sdks/python/test-suites/direct/common.gradle 的batchTests列表中加入测试路径。在 Dataflow Runner 上运行集成测试的步骤构建 SDK tarballcd sdks/python pip install build python -m build --sdisttarball 生成于sdks/python/sdist/目录。通过--test-pipeline-options指定--sdk_location指向 tarball。完整命令示例python -m pytest -o log_cliTrue -o log_levelInfo \ apache_beam/ml/inference/pytorch_inference_it_test.py::PyTorchInference \ --test-pipeline-options--runnerTestDataflowRunner --projectproject --temp_locationgs://bucket/tmp --sdk_locationdist/apache-beam-2.53.0.dev0.tar.gz --regionus-central1如果要让集成测试进入 Python PostCommit 测试套件的 Dataflow 任务使用标记pytest.mark.it_postcommit。为修改后的 SDK 代码构建容器运行./gradlew :sdks:python:container:py39:docker \ -Pdocker-repository-rootgcr.io/location -Pdocker-tagtag推送容器通过--sdk_container_image选项指定容器位置。完整命令示例python -m pytest -o log_cliTrue -o log_levelInfo \ apache_beam/ml/inference/pytorch_inference_it_test.py::PyTorchInference \ --test-pipeline-options--runnerTestDataflowRunner --projectproject --temp_locationgs://bucket/tmp --sdk_container_imageus.gcr.io/apache-beam-testing/beam-sdk/beam:dev --regionus-central1指定额外测试依赖使用--requirements_file选项声明额外依赖python -m pytest -o log_cliTrue -o log_levelInfo \ apache_beam/ml/inference/pytorch_inference_it_test.py::PyTorchInference \ --test-pipeline-options--runnerTestDataflowRunner --projectproject --temp_locationgs://bucket/tmp --sdk_locationus.gcr.io/apache-beam-testing/beam-sdk/beam:dev --regionus-central1 --requirements_filerequirements.txt如果使用的是 Dataflow Runner也可以采用自定义容器方式以官方 Beam SDK 容器镜像为基础镜像在其上应用自己的改动。用修改后的 Beam 代码运行 Pipeline构建 Beam SDK tarball在sdks/python下执行python -m build --sdist详见上节集成测试部分。在 Python 虚拟环境中连同所需扩展安装 SDK命令形如pip install /path/to/apache-beam.tar.gz[gcp]启动你的 Python 脚本命令形如python my_pipeline.py --runnerDataflowRunner --sdk_location/path/to/apache-beam.tar.gz --projectmy_project --regionus-central1 --temp_locationgs://my-bucket/temp ...使用 Dataflow Runner 的两个实用提示Python worker 在处理 work item 前会先安装 Apache Beam SDK因此通常不需要提供自定义 worker 容器。例外是GCP VM 没有外网访问、且临时依赖相对官方发布镜像有变化时才需要自定义 worker 容器做法见上文构建容器一节。从源码安装 Beam Python SDK 可能很慢在n1-standard-1机器上约需 3.5 分钟。替代方案如果宿主机是 amd64 架构可以构建 wheel 而不是 tarball命令形如./gradlew :sdks:python:bdistPy311linux对应 Python 3.11同样通过--sdk_location传入安装只需几秒。注意save_main_session在远程 Runner 上运行DoFn时出现NameError原因是主 Pipeline 模块中的全局 import、函数和变量默认不会被序列化解决方案在 Pipeline 选项中加入--save_main_session。附录格式化 CHANGES.md、常见问题与快照构建目录格式化 CHANGES.md更新CHANGES.md后使用以下 Gradle 命令确保格式正确./gradlew formatChanges该命令会按模板规定顺序组织各章节确保所有必需章节存在在保持既有内容的前提下维持一致的格式。常见问题遇到java.lang.NoClassDefFoundError或 proto 变更相关的诡异错误时依次尝试运行./gradlew clean删除 Gradle 缓存例如rm -fr ~/.gradle与rm -fr beam-repo-dir/.gradle删除仓库根目录下的build目录。用 Gradle 运行单个 Java 测试时用--tests过滤例如./gradlew :it:google-cloud-platform:WordCountIntegrationTest --tests org.apache.beam.it.gcp.WordCountIT.testWordCountDataflow快照构建产物目录位置内容https://repository.apache.org/content/groups/snapshots/org/apache/beam/Java SDK 构建nightlyhttps://gcr.io/apache-beam-testing/beam-sdkBeam SDK 容器构建Java、Python、Go每 4 小时一次https://gcr.io/apache-beam-testing/beam_portabilityportable runnerFlink、Sparkjob server 容器nightlygs://beam-python-nightly-snapshotsPython SDK 构建nightly小结本文围绕contributor-docs/code-change-guide.md的完整脉络梳理了 Beam 代码修改工作的三条主线构建体系Beam 是单一 Gradle 工程BeamModulePlugin通过 natures 机制统一管理 Java/Python/Go/Docker 等模块配置掌握./gradlew :path:task的调用方式即可覆盖绝大多数本地验证需求测试链路Java 侧按**Test.java单元/**IT.java集成区分集成测试通过TestPipeline与beamTestPipelineOptions注入 Runner 和 GCP 参数Python 侧用 pytest 直接驱动--test-pipeline-options控制测试 Pipeline 行为--sdk_location/--sdk_container_image注入本地构建产物改动落地Java 修改经publishToMavenLocal进入 Maven Local配合 SNAPSHOT 仓库配置被用户 Pipeline 拾取Python 修改经 sdist/wheel 构建后由 worker 端安装涉及 harness 或 SDK 容器的修改则需额外构建 shadowJar、harness 容器或自定义 SDK 镜像。所有关键操作均可在仓库内对应文件如 examples/java/build.gradle、sdks/java/io/google-cloud-platform/build.gradle、buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy中查证便于读者按需深入。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 贡献者文档体系从代码修改、测试验证到版本发布的完整工程实践Apache Beam 贡献者文档体系从代码修改、测试验证到版本发布的完整工程实践 Apache Beam 的贡献者文档 contributor docs/大数据批处理流处理数据工程Apache Weex Android SDK 源码构建指南从 Gradle 打包到双包名产物发布Apache Weex Android SDK 源码构建指南从 Gradle 打包到双包名产物发布 导读 本文以 Apache WeexIncubating移动开发跨平台原生移动前端Apache Iceberg源码构建指南从Gradle编译到本地部署Apache Iceberg源码构建指南从Gradle编译到本地部署 引言告别依赖困境掌握源码构建主动权 你是否曾因依赖库版本冲突而陷入困境是否在调试时大数据数据湖OLAP数据存储上一篇FamilyBucket认证授权深度解析JWT无状态认证与动态权限控制实现下一篇Sunshine测试驱动开发Android应用单元测试与集成测试实践指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表