ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam 容器环境(Container Environments)完全指南:构建与运行自定义 SDK 容器镜像

Apache Beam 容器环境(Container Environments)完全指南:构建与运行自定义 SDK 容器镜像 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 通过 SDK Harness 容器SDK container将用户代码的执行环境与 Runner 隔离是支撑 Flink、Spark、Dataflow 等分布式 Runner 在批处理与流处理场景下执行用户代码的核心机制。本文基于官方容器环境指南系统讲解 Beam SDK 容器镜像的发布与拉取方式、三种自定义容器镜像的构建方法基于官方镜像写新 Dockerfile、修改 Beam 源码 Dockerfile、改造现有镜像使其兼容 Beam Runner以及如何在 PortableRunner、FlinkRunner、SparkRunner 和 DataflowRunner 上运行自定义容器镜像并结合仓库源码给出底层实现依据与排错建议。读完本文你将能够为 Beam 流水线定制包含额外依赖、第三方软件或特殊运行环境的 worker 容器并在各 Runner 上正确使用它们。为什么 Beam 需要容器化运行环境Beam 采用 Runner 与 SDK 分离的架构Runner如 Flink、Spark、Dataflow负责调度与执行图而用户代码DoFn、source/sink 等运行在被称为 SDK Harness 的独立进程中。为了让分布在集群各节点的 SDK Harness 拥有一致、可控的运行环境Beam 将每个受支持语言的 SDK 运行时打包为 Docker 容器镜像从而将流水线的执行环境与其他运行时系统隔离开来。这一设计对应的底层契约是 SDK Harness container contract容器内通过/opt/apache/beam/boot这一引导程序boot loader完成与 Runner 的通信、拉取 staged 文件、安装依赖并启动 SDK worker 进程。每个受支持语言在 Beam 发布时都会构建并推送对应的预构建 SDK 容器镜像到 Docker Hub例如apache/beam_python3.12_sdk:2.63.0。需要特别注意的是这些镜像被设计为分布式执行 Beam 流水线的worker 镜像而不是完整的 SDK 开发环境——你可以在其中运行流水线但不应该把它当作日常开发调试用的 Python/Java/Go 环境。从源码看镜像的入口ENTRYPOINT被固定为 boot 引导程序Python SDK 镜像sdks/python/container/Dockerfile以ENTRYPOINT [/opt/apache/beam/boot]结尾Java SDK 镜像sdks/java/container/Dockerfile同样以ENTRYPOINT [/opt/apache/beam/boot]结尾该 boot 程序的实现位于各 SDK 容器目录下的boot.go如 Python 容器 boot.go它解析--id、--logging_endpoint、--artifact_endpoint、--provision_endpoint、--control_endpoint等参数从 Provision service 获取流水线选项与依赖信息下载并安装 staged 的 SDK 源码包、requirements、额外依赖和工作流压缩包最后以python -m apache_beam.runners.worker.sdk_worker_main见 boot.go 中的 sdkHarnessEntrypoint 常量启动真正的 SDK worker。因此任何自定义容器只要保证入口点为/opt/apache/beam/boot或与之等价的自建 boot loader并且镜像内装有与提交环境版本一致的 SDK就能被 Beam Runner 当作合法的 worker 镜像使用。自定义容器的动机与前置条件虽然官方镜像开箱即用但在许多真实场景下你需要对镜像做定制典型动机包括预装额外依赖在镜像中预先安装流水线所需的第三方 Python 包、JAR 或系统库避免每次作业启动时都从远端拉取和安装在 worker 环境中启动第三方软件例如需要在 worker 节点上运行 agent、sidecar 进程或本地服务进一步定制执行环境如更换操作系统基础镜像、调整环境变量、加固安全配置等。前置条件本指南的构建过程依赖 Docker需要先在本地安装 Docker一些 CI/CD 平台如 Google Cloud Build也提供构建镜像的能力。对于远程执行的 Runner你需要一个容器镜像仓库来托管自定义镜像例如 Docker Hub或自建仓库以及云厂商专属仓库如 Google Container Registry (GCR)、Amazon Elastic Container Registry (ECR)。务必确保你的执行引擎/Runner 可以访问该仓库。注意自 2020 年 11 月 20 日起Docker Hub 对匿名和免费认证用户实施了拉取限速rate limits这可能会影响需要多次拉取容器的大规模流水线。建议将自定义镜像托管在 Runner 可访问的私有或云仓库中。为获得最佳体验建议使用最新发布版的 Beam。构建与推送自定义容器三种路径Beam 的 SDK 容器镜像由仓库中检入的 Dockerfile 构建并在每次发布时推送至 Docker Hub。从 sdks/python/container/common.gradle 可以看到Gradle 的docker任务默认将镜像命名为apache/beam_pythonversion_sdk并以gradle.properties中定义的sdk_version作为默认 tag当前仓库的 gradle.properties 中sdk_version2.78.0.dev默认仓库根为apache、前缀为beam_。在此基础上官方文档给出三种定制路径按定制深度递增排列基于已发布镜像编写新的 Dockerfile适合简单扩展如添加文件或环境变量修改 Beam 源码中的 Dockerfile 并从源码构建需要从 Beam 源码构建但允许更深度的定制替换工件、更换基础 OS/语言版本改造现有镜像使其兼容 Beam Runner从已有镜像出发将其配置成可被 Beam Runner 使用的镜像。路径一基于已发布镜像编写新的 Dockerfile这是最简单的方式适合只需在镜像上做少量增量修改的场景。步骤新建 Dockerfile用FROM指令指定基础镜像FROM apache/beam_python3.7_sdk:2.25.0 ENV FOObar COPY /src/path/to/file /dest/path/to/file/此例以预构建的 Python 3.7 SDK 镜像beam_python3.7_sdktag 为 SDK 版本2.25.0为基础额外增加一个环境变量和一个文件。用 Docker 构建并推送镜像export BASE_IMAGEapache/beam_python3.7_sdk:2.25.0 export IMAGE_NAMEmyremoterepo/mybeamsdk # 自定义容器尽量不要使用 latest 标签以便复现失败 export TAGmybeamsdk-versioned-tag # 可选先把基础镜像拉到本地 Docker daemon确保本地是最新版 docker pull ${BASE_IMAGE} docker build -f Dockerfile -t ${IMAGE_NAME}:${TAG} .如果 Runner 远程运行需要重新打 tag 并推送镜像到合适的仓库docker push ${IMAGE_NAME}:${TAG}推送完成后核对远程镜像的 ID 和 digest 与docker build或docker images输出的本地 ID、digest 一致确保推送成功。路径二修改 Beam 源码中的 Dockerfile 并从源码构建此方式需要从 Beam 源码构建镜像工件适合需要替换基础 OS/语言版本、替换内置工件等深度定制的场景。开发环境的搭建可参考 Contribution guide。注意建议从与你运行流水线所用 SDK 版本一致的稳定发布分支release-X.XX.X开始。SDK 版本不一致可能导致意想不到的错误。具体步骤克隆 beam 仓库并切换到对应发布分支export BEAM_SDK_VERSION2.26.0 git clone https://github.com/apache/beam.git cd beam # 保存当前目录为工作目录 export BEAM_WORKDIR$PWD git checkout origin/release-$BEAM_SDK_VERSION定制对应语言的 Dockerfile通常位于sdks/language/container/Dockerfile例如 Python 的 Dockerfile。以 Python 为例可以修改基础镜像FROM python:${py_version}-bookworm、增减apt包、调整pip install的依赖列表依赖清单来自base_image_requirements.txt由 common.gradle 中的 generatePythonRequirements 任务 根据[gcp,dataframe,test,tfrecord,yaml]等 extras 生成Java 镜像则对应 sdks/java/container/Dockerfile其基础镜像由base_image与java_version两个 ARG 决定。回到 Beam 根目录运行对应镜像的 Gradledocker目标cd $BEAM_WORKDIR # 各 SDK 的默认仓库 ./gradlew :sdks:java:container:java11:docker ./gradlew :sdks:java:container:java17:docker ./gradlew :sdks:java:container:java21:docker ./gradlew :sdks:go:container:docker ./gradlew :sdks:python:container:py310:docker ./gradlew :sdks:python:container:py311:docker ./gradlew :sdks:python:container:py312:docker ./gradlew :sdks:python:container:py313:docker # 构建所有 Python SDK 镜像的快捷方式 ./gradlew :sdks:python:container:buildAll从仓库目录结构看当前仓库支持 Java 11/17/21sdks/java/container/java11、java17、java21 等子目录、Python 3.103.14sdks/python/container/py310 至 py314以及 Go 的容器构建。Python 的完整构建流程还可参考 Python SDK 镜像构建指南该文档详细说明了本地环境要求Java、Python、Golang、Docker 与 docker-buildx、./gradlew :sdks:python:container:py版本:docker的使用方式以及将镜像推送至 Artifact Registry / Docker Hub 的命令。用docker images验证镜像已生成$ docker images --digests REPOSITORY TAG DIGEST IMAGE ID CREATED SIZE apache/beam_java11_sdk latest sha256:... ... 1 min ago ... apache/beam_java17_sdk latest sha256:... ... 1 min ago ... apache/beam_java21_sdk latest sha256:... ... 1 min ago ... apache/beam_python3.10_sdk latest sha256:... ... 1 min ago ... apache/beam_python3.11_sdk latest sha256:... ... 1 min ago ... apache/beam_python3.12_sdk latest sha256:... ... 1 min ago ... apache/beam_python3.13_sdk latest sha256:... ... 1 min ago ... apache/beam_go_sdk latest sha256:... ... 1 min ago ...若 Runner 远程运行重新打 tag 并推送若通过附加构建参数指定了自定义 repo/tag 则可跳过此步export BEAM_SDK_VERSION2.26.0 export IMAGE_NAMEgcr.io/my-gcp-project/beam_python3.7_sdk export TAG${BEAM_SDK_VERSION}-custom docker tag apache/beam_python3.7_sdk ${IMAGE_NAME}:${TAG} docker push ${IMAGE_NAME}:${TAG}推送后核对远程与本地镜像 ID/digest 是否一致。附加构建参数Gradle 的 docker 任务为镜像定义了默认仓库与 tag默认仓库是 Docker Hub 的apache命名空间默认 tag 是 gradle.properties 中定义的 SDK 版本。可以通过给构建任务传参来指定不同的仓库或 tag例如./gradlew :sdks:python:container:py36:docker -Pdocker-repository-rootexample-repo -Pdocker-tag2.26.0-custom该命令构建 Python 3.6 容器并打 tag 为example-repo/beam_python3.6_sdk:2.26.0-custom。从 Beam 2.21.0 起引入了docker-pull-licenses标志用于把第三方依赖的许可证/声明文件加入镜像。例如./gradlew :sdks:java:container:java11:docker -Pdocker-pull-licenses会构建一个带许可证的 Java 11 SDK 镜像许可证存放在/opt/apache/beam/third_party_licenses/。默认情况下镜像不包含这些许可证文件。从源码可以印证这一点Python 容器 common.gradle 将pull_licenses作为 build arg 传入且仅在指定了docker-pull-licenses属性或是 release 构建时才为 truePython Dockerfile 中则用多阶段构建third_party_licenses阶段收集许可证并在非 pull 模式下删除/opt/apache/beam/third_party_licenses见 Python Dockerfile 第 110-128 行Java Dockerfile 也有同样的pull_licenses判断逻辑sdks/java/container/Dockerfile#L52-L55。路径三改造现有镜像使其兼容 Beam Runner如果你不想从 Beam 镜像出发而希望以自己已有的镜像为基础构建自定义容器最省事的方式是使用 Docker多阶段构建multi-stage build从默认的 Apache Beam 基础镜像中把必要工件复制到你的自定义镜像里。将必要工件从 Apache Beam 基础镜像复制到你的镜像# 这可以是任意容器镜像 FROM python:3.8-bookworm # 安装 SDKPython SDK 需要 RUN pip install --no-cache-dir apache-beam[gcp]2.52.0 # 从官方 SDK 镜像复制文件包括脚本与依赖 COPY --fromapache/beam_python3.8_sdk:2.52.0 /opt/apache/beam /opt/apache/beam # 按需做其他定制 # 将入口点设置为 Apache Beam SDK launcher ENTRYPOINT [/opt/apache/beam/boot]注意本例假定现有基础镜像已装好所需依赖此处是 Python 3.8 与 pip。把 Apache Beam SDK 装进镜像可以确保镜像包含必要的 SDK 依赖并缩短 worker 启动时间。RUN指令中指定的版本必须与启动流水线所用的版本一致。务必保证基础镜像中的 Python 或 Java 运行时版本与运行流水线所用的版本相同。注意自定义镜像中任何额外的 Python 依赖都应安装到全局 Python 环境中。这是因为 Python 容器 boot.go 会为每个 worker 创建独立的虚拟环境venv默认位于/opt/apache/beam-venv但它以--system-site-packages方式创建并继承全局包同时 boot.go 在installSetupPackages中会校验apache-beam是否已安装——若未安装会直接报错并提示“如果你使用自定义容器镜像必须在镜像中安装与流水线提交环境相同版本的 apache-beam 包”见 boot.go 第 464-468 行。依赖装到全局环境可以确保新创建的 venv 能访问到它们。用 Docker 构建并推送export BASE_IMAGEapache/beam_python3.8_sdk:2.52.0 export IMAGE_NAMEmyremoterepo/mybeamsdk export TAGlatest # 可选先拉取基础镜像确保本地是最新版 docker pull ${BASE_IMAGE} docker build -f Dockerfile -t ${IMAGE_NAME}:${TAG} .若 Runner 远程运行重新打 tag 并推送docker push ${IMAGE_NAME}:${TAG}从零构建兼容容器Go SDK 专属自 Beam 2.55.0 起Go SDK 改用 distroless 镜像 作为基础镜像。这类镜像通过不包含常见工具和实用程序来减小安全攻击面但这也会让前面介绍的定制方式变得困难。作为后备方案你可以从零构建自定义镜像构建一个配套的 boot loader并把它设置为容器入口点。例如如果你想用 alpine 作为容器操作系统多阶段 Dockerfile 大致如下FROM golang:latest-alpine AS build_base # 设置容器内当前工作目录 WORKDIR /tmp/beam # 构建与你 Beam 版本匹配的 Beam Go bootloader 到本地目录。 # 其他 SDK 语言也有类似的 go target。 RUN GOBINpwd go install github.com/apache/beam/sdks/v2/go/containerv2.53.0 # 设置真正的基础镜像 FROM alpine:3.9 RUN apk add ca-certificates # 以下内容为容器正确运行所必需 # 将 boot loader container 复制到镜像中 COPY --frombuild_base /tmp/beam/container /opt/apache/beam/boot # 设置容器使用新构建的 boot loader ENTRYPOINT [/opt/apache/beam/boot]构建与推送方式同改造现有基础镜像一致。注意Java 和 Python 还需要额外的依赖如各自的语言运行时和 SDK 包才能构成有效容器镜像对这两种 SDK 而言仅靠 boot loader 不足以创建自定义容器。使用自定义容器镜像运行流水线构建好自定义镜像后需要让 Runner 在启动流水线时使用它。通用的方式是通过 PortableRunner 的--environment_config参数或支持 PortableRunner 标志的其他 Runner指定容器镜像其他 Runner如 Dataflow则使用不同的参数。--environment_type与--environment_config的定义可参见 PortablePipelineOptionsenvironment_type的可选值为DOCKER和PROCESSenvironment_config在DOCKER类型下为 Docker 镜像 URL在PROCESS类型下为 JSON 形式的进程配置。Python 侧对应的解析逻辑位于 portable_runner.py。在 PortableRunner 上运行export IMAGEmy-repo/beam_python_sdk_custom export TAGX.Y.Z export IMAGE_URL${IMAGE}:${TAG} python -m apache_beam.examples.wordcount \ --input/path/to/inputfile \ --output /path/to/write/counts \ --runnerPortableRunner \ --job_endpointembed \ --environment_typeDOCKER \ --environment_config${IMAGE_URL}在 FlinkRunner 上运行export IMAGEmy-repo/beam_python_sdk_custom export TAGX.Y.Z export IMAGE_URL${IMAGE}:${TAG} # 使用 FlinkRunner 运行流水线它会启动一个 Flink job server python -m apache_beam.examples.wordcount \ --input/path/to/inputfile \ --outputpath/to/write/counts \ --runnerFlinkRunner \ # 本地运行批作业时需要复用容器 --environment_cache_millis10000 \ --environment_typeDOCKER \ --environment_config${IMAGE_URL}在 SparkRunner 上运行export IMAGEmy-repo/beam_python_sdk_custom export TAGX.Y.Z export IMAGE_URL${IMAGE}:${TAG} # 使用 SparkRunner 运行流水线它会启动 Spark job server python -m apache_beam.examples.wordcount \ --input/path/to/inputfile \ --outputpath/to/write/counts \ --runnerSparkRunner \ # 本地运行批作业时需要复用容器 --environment_cache_millis10000 \ --environment_typeDOCKER \ --environment_config${IMAGE_URL}--environment_cache_millis的语义可以在 PortablePipelineOptions 中找到它表示作业内环境缓存的毫秒数0 表示不缓存。本地批处理场景下设置一个非零值可以复用同一个容器执行多个 bundle显著降低启动开销。在 DataflowRunner 上运行Dataflow 使用不同的参数--sdk_container_image指定容器镜像export GCS_PATHgs://my-gcs-bucket export GCP_PROJECTmy-gcp-project export REGIONus-central1 # 默认情况下Dataflow Runner 可以访问同一项目下的 GCR 镜像 export IMAGEmy-repo/beam_python_sdk_custom export TAGX.Y.Z export IMAGE_URL${IMAGE}:${TAG} # 在 Dataflow 上运行流水线。 # 这是 Python 批处理流水线要在 Dataflow Runner V2 上运行 # 必须指定实验 use_runner_v2 python -m apache_beam.examples.wordcount \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output ${GCS_PATH}/counts \ --runner DataflowRunner \ --project $GCP_PROJECT \ --region $REGION \ --temp_location ${GCS_PATH}/tmp/ \ --experimentuse_runner_v2 \ --sdk_container_image$IMAGE_URL关于镜像 tag 的实践建议避免对自定义镜像使用:latest标签。请用日期或唯一标识符给构建打 tag。一旦出问题这种有意义的 tag 能让你把流水线回退到之前已知可用的配置并便于检查变更内容。这一点也与 sdk_container_builder.py 中“基于latest或固定版本构造基础镜像名”的逻辑形成呼应——可复现性始终是自定义容器实践的关键。进阶利用 boot 的 setup-only 模式预构建带依赖的镜像除了手工编写 DockerfileBeam Python SDK 还提供了一种“先把依赖装进镜像”的自动化途径SdkContainerImageBuilder实现于 sdk_container_builder.py会把 Beam SDK 源码包、requirements.txt依赖、extra_packages.txt中的包以及 workflow 压缩包复制到最新的官方 Python SDK 容器镜像中然后以--setup_only --artifacts manifest的方式提前运行 boot 程序完成依赖安装从而构建出新的镜像对应的 Dockerfile 模板见 sdk_container_builder.py 第 55-60 行。boot.go 中setup_only模式的实现在 processArtifactsInSetupOnlyMode它只处理工件并安装依赖然后直接退出不启动 worker 进程——这样产出的镜像在流水线运行时无需再安装依赖可显著缩短 worker 启动时间。这一能力与上述“在镜像中预装额外依赖”的动机完全一致适合对启动延迟敏感的生产流水线。故障排查Troubleshooting使用自定义容器运行 Beam 流水线遇到意外错误时可以从以下几个方面排查语言与 SDK 版本不一致容器内 SDK 与流水线 SDK 的语言/版本差异可能因不兼容导致意外错误。最佳实践是基础容器与运行流水线时使用相同的稳定 SDK 版本。这一点在源码中也有体现——Python boot.go 会显式校验apache-beam包已安装并提示版本须与提交环境一致。远程镜像不可访问使用远程容器遇到意外错误时确认镜像存在于远程仓库且如有必要任何第三方服务都能访问该仓库。本地拉取失败本地 Runner 会先尝试拉取远程镜像拉不到时回退到本地镜像。如果镜像无法被本地 Docker daemon 拉取可能会看到类似日志Error response from daemon: manifest for remote.repo/beam_python3.7_sdk:2.25.0-custom not found: manifest unknown: ... INFO:apache_beam.runners.portability.fn_api_runner.worker_handlers:Unable to pull image...出现这类信息时应核对镜像 tag 是否拼写正确、镜像是否已 push 到目标仓库以及本地/远程的镜像 ID 与 digest 是否一致。小结Beam 容器环境机制的核心是“Runner 与 SDK Harness 分离”的架构每个语言的 SDK 运行时被打包为以/opt/apache/beam/boot为入口的容器镜像由 Runner 按需拉起执行用户代码。自定义容器的三条路径覆盖了从“简单加依赖”到“彻底换基础镜像/OS”的全部需求基于官方镜像写 Dockerfile路径一、修改 Beam 源码 Dockerfile 并用 Gradle 构建路径二、多阶段构建从官方镜像复制工件路径三以及 Go SDK 特有的从零构建方案。运行侧只需通过--environment_type/--environment_configPortable 系 Runner或--sdk_container_imageDataflow指定镜像即可。始终记住两条铁律容器内 SDK/语言运行时版本必须与流水线提交环境一致自定义镜像避免使用:latest标签——前者保证兼容性后者保证可复现性。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam SDK 自定义容器镜像构建指南基于 Python/Java SDK 的三种定制方法Apache Beam SDK 自定义容器镜像构建指南基于 Python/Java SDK 的三种定制方法 Apache Beam 的 SDK 运行环境可以用批处理流处理大数据TileLang 容器化环境搭建Docker 镜像构建与 GPU 容器运行全解TileLang 容器化环境搭建Docker 镜像构建与 GPU 容器运行全解 TileLang 的构建依赖一条较长的工具链定制版 TVM 子模块、CUTL编译器编程语言高性能计算人工智能深度学习为 Apache Beam SDK 定制容器环境三种 Docker 镜像定制方法全解析为 Apache Beam SDK 定制容器环境三种 Docker 镜像定制方法全解析 本篇技术指南围绕 Apache Beam 官方文档《Container大数据批处理流处理数据工程上一篇Robomongo 内嵌 esprima 2.7.3ECMAScript 解析器在 MongoDB Shell 脚本解析中的集成与应用下一篇bb 开发指南如何参与 bb 单体仓库开发pnpm dev、Turbo、worktree 完整流程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进