多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Beam Python 流水线依赖管理完全指南:从 requirements.txt 到自定义容器与可复现环境

Apache Beam Python 流水线依赖管理完全指南:从 requirements.txt 到自定义容器与可复现环境 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载依赖管理是 Apache Beam Python 流水线从本地能跑走向远程可靠运行的关键一环。本文以官方文档《Managing Python Pipeline Dependencies》为核心骨架结合仓库中SetupOptions参数体系与 Juliaset 示例等源码证据系统讲解 PyPI 依赖的分发、本地与非公开包的上传、多文件工程的打包、非 Python 依赖的安装、SDK 容器预构建以及启动环境与运行时环境的一致性与可复现性。读完本文你将掌握从requirements.txt到setup.py、--extra_package、自定义容器等一整套依赖供给与管控方案并能诊断远程执行中的NameError、ModuleNotFound等典型依赖问题。依赖管理的本质指定依赖并控制在生产环境中的使用依赖管理要解决两件事指定你的流水线需要哪些依赖以及控制生产环境实际使用哪些依赖。理解这一点前需要先明确运行模型流水线代码在**提交端launch environment**完成图构建与序列化真正执行数据处理的则是在远程workerruntime environment上。远程 worker 通常运行在基于 Debian 的容器镜像中自带一套标准的 Python 发行版。前置判断如果你的代码只依赖 Python 标准库那么远程 worker 的默认环境已经够用本页介绍的所有操作都可以跳过——你不需要做任何事。只有当代码依赖了标准库之外的包PyPI 上的第三方包、非公开的本地包、甚至需要apt安装的系统库时才需要把依赖搬运到 worker 上。Beam Python SDK 为此提供了一条完整的依赖供给链路提交时下载/打包 → 暂存staging到 runner → 运行时在 worker 上安装。这条链路的入口就是SetupOptions定义的一系列 pipeline options。PyPI 依赖requirements.txt与--requirements_file如果流水线使用 Python Package Index 上的公开包就必须让这些包在远程 worker 上可用。对于单个 Python 文件或 notebook 组成的简单流水线最直接的方式是提供一个requirements.txt文件对于更复杂的场景则应把流水线组织成包并考虑用自定义容器安装依赖。生成与整理 requirements.txt找出本机安装的包执行pip freeze requirements.txt该命令会列出本机安装的所有包无论来自哪个安装源并写入requirements.txt。裁剪编辑该文件删除与你的流水线代码无关的包。pip freeze是全量快照会包含大量无关的传递依赖直接全量提交既臃肿又可能引入版本冲突。提交时安装运行流水线时带上--requirements_file requirements.txtrunner 会根据该文件把额外依赖安装到远程 worker 上。提示除了pip freeze也可以使用 pip-tools 之类的工具在requirements.in中只写顶层依赖由工具编译出完整依赖清单这样生成的 requirements 文件更干净、更可控。依赖暂存requirements cache机制--requirements_file背后有一套下载-缓存-暂存机制仓库源码中的帮助文本给出了权威描述pipeline_options.py在流水线提交阶段Beam 会先把 requirements 文件中指定的包下载到本地的 requirements cache 目录随后把这个 cache 目录暂存stage到 runner在运行时只要 cache 可用worker 就从 cache 安装包无需再连接 PyPIcache 按需刷新已有包不会重复下载。这意味着只要提交时能访问 PyPI运行时就未必需要——对运行在受限网络环境中的流水线非常有用。相关选项--requirements_cache 目录修改缓存目录的默认位置--requirements_cacheskip跳过 requirements 暂存此时 worker 运行时会直接访问 PyPI--requirements_cache_only_sources只缓存源码分发包sdist。源码注释特别标注了此选项可能显著拖慢提交且是为保留 2.37.0 之前的行为而设未来版本可能移除pipeline_options.py。自定义容器--sdk_container_image如果依赖较多或安装耗时较长可以把所有依赖预先烧进容器镜像运行时直接用该镜像启动 worker构建自定义容器镜像在构建时把--requirements_file中的依赖装进镜像# 将以下行加入 Dockerfile并把 path to requirements.txt 换成实际路径 COPY path to requirements.txt /tmp/requirements.txt RUN python -m pip install -r /tmp/requirements.txt运行流水线时通过--sdk_container_image image指定镜像。收益依赖在镜像构建阶段就安装完成运行时不再需要传--requirements_file从而显著缩短 worker 启动时间省去每次运行时在线安装的开销。本地 Python 包或非公开依赖--extra_package如果依赖无法从公开 PyPI 获取例如从内部仓库下载的包、闭源或未发布的包需要显式把包文件传给 runner识别非公开包pip freeze运行流水线时指定包文件--extra_package /path/to/package/package-name其中package-name是包的tarball。可以用build工具构建pip install --upgrade build python -m build --sdist关于--extra_package的格式要求源码给出的权威定义是pipeline_options.py文件必须是 (1) 包 tarball.tar、(2) 压缩包 tarball.tar.gz、(3) Wheel 文件.whl或 (4) 可用标准pip install安装的压缩 zip.zip。多个包可以重复传--extra_packageworker 会按命令行顺序安装文件在提交时被暂存到--staging_location指定的暂存区。暂存单个文件--files_to_stage与 worker 端访问如果流水线只依赖一个或几个独立的 Python 文件、或非 Python 数据文件JSON 配置、模型文件等没必要打包成完整 Python 包可以逐个暂存--files_to_stage/path/to/my_module.py,/path/to/data_config.json--files_to_stage别名--file_to_stage接受逗号分隔的本地文件路径列表。运行时 runner 上传这些文件并在 worker 的 staged files 目录中提供它们pipeline_options.py。worker 上如何访问暂存文件Python 模块从 Apache Beam 2.76.0 起SDK worker 启动时会把 staged files 目录自动追加到sys.path因此可以直接在代码里import my_module数据/配置文件非 Python 文件会被下载到 staged files 目录。worker 上通过SEMI_PERSISTENT_DIRECTORY环境变量定位该目录容器化 runner 上默认是/tmp/stagedimport os staged_dir os.environ.get(SEMI_PERSISTENT_DIRECTORY, /tmp/staged) config_path os.path.join(staged_dir, data_config.json)与--beam_plugins组合worker 启动即导入从 Apache Beam 2.76.0 起可以把暂存文件与--beam_plugins配合让 worker 在任何处理开始前先执行初始化代码用--file_to_stage暂存插件文件如my_custom_plugin.py在--beam_plugins中引用其模块名--file_to_stage/path/to/my_custom_plugin.py \ --beam_pluginsmy_custom_plugin这会指示 SDK worker 在启动时立即import my_custom_plugin从而触发模块中的初始化逻辑。从源码看该机制在 worker 侧有明确实现SDK worker 启动流程中的_import_beam_plugins函数会依次导入列表中的插件模块sdk_worker_main.py插件列表来自SetupOptions.beam_pluginssdk_worker_main.pyDataflow runner 侧还会把自动收集的插件与用户指定的插件合并去重dataflow_runner.py。仓库中还有专门的集成测试验证插件暂存与导入行为beam_plugins_it_test.py。注意源码帮助文本注明--beam_plugins目前仍是实验性选项不提供稳定性保证pipeline_options.py。多文件依赖setup.py与--setup_file流水线代码通常横跨多个文件。对远程运行来说最稳妥的做法是把这些文件组织成 Python 包提交时把包暂存worker 启动时自动安装。完整操作步骤创建setup.py一个最简模板import setuptools setuptools.setup( namePACKAGE-NAME, versionPACKAGE-VERSION, install_requires[ # 列出流水线依赖的 Python 包 ], packagessetuptools.find_packages(), )组织项目结构根目录包含setup.py、主工作流文件和一个放其余代码的包目录root_dir/ setup.py main.py my_package/ __init__.py my_pipeline_launcher.py my_custom_dofns_and_transforms.py other_utils_and_helpers.py仓库中的Juliaset示例正是这种结构的完整范例根目录是 juliaset_main.py 与 setup.py包代码位于juliaset/子目录juliaset/juliaset.py并配有 单元测试。在提交环境安装包本地开发环境与 worker 保持一致的入口pip install -e .运行流水线时指定 setup 文件--setup_file /path/to/setup.py源码对--setup_file的行为定义是提交时构建源码分发包sdistworker 在运行任何自定义代码之前先安装该包pipeline_options.py。与--requirements_file的关键区别若包的依赖已在setup.py的install_requires字段中声明无需再传--requirements_file但使用--setup_file时Beam 不会把依赖包暂存到 runner只有流水线包本身被暂存如果运行环境没有这些依赖worker 会在运行时从 PyPI 安装。这正是不使用自定义容器时可复现运行时环境需要额外手段见下文的原因。setuptools 兼容性陷阱使用install_requires时要确保所列包与 runner 环境所用的 setuptools 版本兼容。如果某些依赖需要向后兼容的旧版 setuptools才能构建成功可以用setup_requires强制指定版本需求。例如为兼容仍依赖现已弃用的pkg_resources模块的包import setuptools setuptools.setup( namePACKAGE-NAME, versionPACKAGE-VERSION, setup_requires[setuptools82.0.0], install_requires[incompatible-package, ...], packagessetuptools.find_packages() )非 Python 依赖用CUSTOM_COMMANDS在 worker 上执行安装命令如果流水线依赖非 Python 包如需要apt install的系统库或者某个 PyPI 包在安装阶段依赖非 Python 组件首选方案是使用自定义容器在构建镜像时一并安装。不使用容器时则必须按上文把流水线组织成包在setup.py的CUSTOM_COMMANDS列表中加入所需的安装命令例如CUSTOM_COMMANDS [ [apt-get, update], [apt-get, --assume-yes, install, libjpeg62], ]仓库中的 Juliaset 示例对这套机制有完整实现定义了CustomCommands这个setuptools.Command子类用subprocess逐条执行CUSTOM_COMMANDS非零退出码即视为失败抛出异常再通过cmdclass把自定义build命令挂到构建流程上确保pip install包时自动触发juliaset/setup.py。源码注释还给出了两条重要实践使用apt-get时第一跳必须是apt-get update否则后续安装会报package not found--assume-yes可跳过交互确认若某个 PyPI 包依赖自定义命令安装的库应把该包也挪进自定义命令列表如[pip, install, my_package]因为自定义命令在 pip 安装依赖之后执行。运行流水线--setup_file /path/to/setup.py必须验证这些命令能在远程 worker 上执行例如用apt的前提是远程 worker 具备apt支持。另外注意因为自定义命令在依赖安装pip之后才执行所以不要把此类依赖同时写进requirements.txt或setup.py的install_requires中避免安装顺序冲突。预构建 SDK 容器镜像--prebuild_sdk_container_engine在 runner 以 Docker 容器启动 SDK worker 的执行模式下--requirements_file等选项指定的附加依赖是在运行时装进容器的这会拉长 worker 启动时间。Beam 提供了在 worker 启动前一次性完成依赖安装的预构建机制--prebuild_sdk_container_engine源码帮助文本定义了引擎的两种取值pipeline_options.pylocal_docker使用本机 Docker 构建cloud_build使用 Google Cloud Build需要已启用 Cloud Build API 的 GCP 项目。启用后SDK 会在提交前于 SDK worker 容器中执行 boot 序列、把全部流水线依赖装进容器并以预构建镜像运行流水线也可以继承SdkContainerImageBuilder在其他环境构建。相关辅助选项包括--cloud_build_machine_type指定 Cloud Build 机型与--docker_registry_push_url预构建镜像的推送仓库pipeline_options.py。限制该特性仅对 Dataflow Runner v2 可用。Google Cloud Dataflow 的具体用法含在自定义镜像中预装额外依赖的指引可参见 Dataflow 官方文档关于 custom containers 的 Pre-building 章节。Pickling 与主会话管理远程 runner 执行时流水线内容如 transform 的用户代码会被序列化pickle成字节码。序列化发生在作业提交阶段反序列化发生在运行时因此提交端与运行端必须使用同一版本的 pickling 库。默认 pickler 的演进Apache Beam 2.64.0 及更早默认 pickler 是dill。使用dill时主流水线模块中定义的全局导入、函数、变量默认不会随作业序列化保存远程 runner 上执行DoFn时可能触发意外的NameError。Apache Beam 2.65.0 及之后默认 pickler 切换为cloudpickle。仓库源码中可见DEFAULT_PICKLE_LIB USE_CLOUDPICKLE且cloudpickle_pickler被设为期望的 picklerpickler.py。--save_main_session使用dill时解决NameError的办法是把主会话内容随流水线提交--save_main_session该选项会把全局命名空间的 pickled 状态加载到 worker 上以DataflowRunner为例。而使用cloudpickle时不需要--save_main_session默认已开启。在旧版 Beam 上想改用 cloudpickle可传--pickle_librarycloudpickle源码确认了默认值的联动逻辑对 service runnerpickle_library 为cloudpickle时save_main_session自动置为Truedill时置为Falsepipeline_options.py。注意--pickle_library的可选值还包括dill_unsafe——这是为绕过版本校验而设的不保证安全选项pipeline_options.py。版本一致性要求使用dill且安装了自定义版本的 dill 的用户必须保证运行时自定义容器或依赖清单中使用同一版本Beam 对序列化库dill、cloudpickle的版本设置了严格约束validator 明确要求dill0.3.1.1版本不符会报错除非使用dill_unsafepipeline_options_validator.py。这一约束正是为了保证提交端与运行端的序列化格式一致。控制流水线使用的依赖构建可复现的环境两个环境的职责划分启动环境launch environment把流水线翻译为 runner 无关表示runner-independent representation并提交执行的环境。可以从安装了 Beam SDK 的 Python 虚拟环境启动也可以用 Dataflow Flex Templates、Notebook 环境、Apache Airflow 等工具启动。运行时环境runtime environmentrunner 执行流水线时使用的 Python 环境流水线代码在此完成数据处理包含 Apache Beam 本身与流水线运行时依赖。构建可复现环境的四种工具requirements 文件安装依赖后pip freeze requirements.txt重建环境时pip install -r requirements.txtconstraint 文件限制包的安装只允许指定版本lock 文件用 PipEnv、Poetry、pip-tools 等工具声明顶层依赖生成钉死版本的全部传递依赖 lock 文件并从 lockfile 创建虚拟环境Docker 容器镜像把启动与运行时环境一并打包进镜像环境只在镜像重建时变化。同时用版本控制管理这些定义环境的配置文件。让运行时环境可复现远程 runner 上使用可复现的运行时环境意味着 worker 每次运行使用相同的依赖不受直接或传递依赖新版本发布的副作用影响也不需要在运行时做依赖解析。两种实现方式使用包含全部依赖的自定义容器镜像通过--sdk_container_image指定在--requirements_file中给出穷尽的依赖清单并用--prebuild_sdk_container_engine在流水线执行前完成运行时环境初始化若依赖不变可用--sdk_container_image复用预构建镜像。自检方法自包含的运行时环境通常是可复现的——可以尝试限制运行时对 PyPI 的互联网访问来验证。若使用 Dataflow Runner可参考其--no_use_public_ips选项关闭外部 IP。受控地重建或升级运行时环境不要修改正在被运行中流水线使用的容器镜像自定义镜像避免使用:latest标签改用日期或唯一标识符打标签便于出问题时回滚到已知可用配置并审查变更考虑把pip freeze输出或requirements.txt内容纳入版本控制系统。让启动环境可复现启动环境运行生产版本的流水线本地开发时你可能另有一个包含 Jupyter、Pylint 等开发依赖的开发环境而生产启动环境未必需要这些额外依赖应单独构建维护。减少提交副作用的关键是能按可复现方式重建启动环境Dataflow Flex Templates 就是容器化、可复现启动环境的实例。要把 Beam 安装到干净的虚拟环境中且做到可复现可以使用 requirements 文件 约束文件约束内容来自 Beam 默认容器镜像自带的依赖清单。命令形如BEAM_VERSION2.48.0 PYTHON_VERSION$(python -c import sys; print(f{sys.version_info.major}{sys.version_info.minor})) pip install apache-beam$BEAM_VERSION \ --constraint /path/to/sdks/python/container/py${PYTHON_VERSION}/base_image_requirements.txt仓库中这些按 Python 版本划分的约束文件真实存在sdks/python/container/py310/base_image_requirements.txt至py314/base_image_requirements.txt另有ml/py3xx下的 ML 镜像版本。约束文件可确保启动环境中 Beam 的依赖版本与默认 Beam 容器一致还可能省去安装时的依赖解析。让启动环境与运行时环境兼容启动环境把流水线图翻译为 runner 无关表示涉及序列化 transform 代码worker 上会反序列化这些内容。如果运行时环境与启动环境差异过大可能因以下原因报错Beam 版本必须一致Python 主/次版本major.minor也必须一致否则可能报Pipeline construction environment and pipeline runtime environment are not compatible旧版 SDK 上可能报SystemError: unknown opcodeprotobuf 版本需要匹配或兼容流水线代码用到的库需要匹配若序列化代码引用了 worker 上不存在的函数或模块远程 runner 会抛ModuleNotFound或AttributeError。遇到这类错误时确认相关库已装到 worker 上并检查是否需要保存主会话pickling 库版本必须在提交与运行时一致如前述dill0.3.1.1的强制约束。若确实要强制安装不同版本的dill/cloudpickle需满足两个条件提交端与运行端安装同一版本且该版本在你的流水线上可用。排查手段对比两个环境的pip freeze输出检查差异升级到最新版 Beam因为新版本 SDK 内置了更多环境兼容性检查。终极方案直接从运行时使用的同一个容器化环境发起流水线——Dataflow 支持用自定义容器镜像构建 Flex Templates实现启动环境 运行时环境。由于两个容器由同一镜像创建启动与运行时环境天然兼容同时两个环境都可以可复现地重建。总结与决策速查场景首选方案相关 pipeline option单文件/notebook依赖公开 PyPI 包requirements.txt--requirements_file依赖多、安装慢、网络受限自定义容器预装依赖--sdk_container_image非公开/本地包tarball/whl/zip提交包文件--extra_package单个 Python/数据文件逐个暂存--files_to_stage/--file_to_stage多文件工程打包为 Python 包--setup_file需要apt等非 Python 依赖自定义容器或CUSTOM_COMMANDS容器 /--setup_file缩短容器启动时间预构建 SDK 镜像--prebuild_sdk_container_engine仅 Dataflow Runner v2dill下远程NameError提交主会话--save_main_session或用--pickle_librarycloudpickle环境可复现与兼容约束文件 / 容器化启动 运行--requirements_file--sdk_container_image等实践要点可以浓缩为三条提交端与运行端版本一致Beam、Python、protobuf、pickle 库、能离线自包含就自包含cache、预构建镜像、自定义容器、所有环境定义文件进版本控制。遵循这些原则远程流水线的依赖问题大多可以在提交阶段就被发现而不是等到 worker 上运行到一半才抛异常。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Kedro 项目依赖管理实战从 requirements.txt 声明到可复现锁定环境Kedro 项目依赖管理实战从 requirements.txt 声明到可复现锁定环境 导读 Kedro 将依赖管理划分为核心模块与项目专属依赖两个层数据工程工作流自动化抖音无水印批量下载实战3步搭好你的个人素材库抖音无水印批量下载实战3步搭好你的个人素材库 手动复制链接一个个存文件散在各处发布时间和评论数据全丢——这是做抖音内容备份时的真实处境。douyin do大数据批处理流处理数据工程爱享素材下载器网络资源嗅探下载与视频号保存教程爱享素材下载器网络资源嗅探下载与视频号保存教程 平台上的视频、图片、音乐想存下来却到处找不到下载入口爱享素材下载器res downloader用本机代桌面应用网络音视频上一篇3分钟QQ音乐解密指南快速解锁加密音乐文件的终极方案下一篇终极解放ncmdump 实现网易云NCM音乐格式自动化解密转换创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表