news 2026/10/3 2:22:35

Cloud Composer DAG 批量暂停/恢复脚本全解析:composer_dags.py 使用与源码解读

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Cloud Composer DAG 批量暂停/恢复脚本全解析:composer_dags.py 使用与源码解读
  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

导读

本文围绕 python-docs-samples 仓库中 Composer DAGs Pausing/Unpausing script 文档 展开,系统讲解用于 Cloud Composer 环境的 DAG 批量管理工具 composer_dags.py。你将掌握该脚本的前置条件、命令参数、pause/unpause 两种操作的完整用法,以及从源码层面理解它如何自动识别 Airflow 版本、枚举环境内全部 DAG 并逐一批量执行操作,从而在迁移、维护、停服等场景下高效管理整个 Composer 环境。

脚本定位:一站式批量管理 Composer 环境中的全部 DAG

在 Cloud Composer 环境中,Airflow 的 Web UI 通常只支持逐个暂停或恢复 DAG。当环境中存在大量 DAG(例如停机维护、版本迁移或灾备演练时),逐个操作既耗时又易遗漏。composer_dags.py正是为解决这一问题而生的独立命令行脚本:它以gcloud composer environments run为底层调用链,自动列举指定 Composer 环境中的全部 DAG,并统一执行暂停(pause)或恢复(unpause)操作。

脚本的核心设计目标(源自 composer_dags.md 与 源码 docstring)是:

  • 批量暂停:将环境中除airflow_monitoring之外的所有 DAG 全部暂停;
  • 批量恢复:将环境中所有 DAG 全部恢复(unpause);
  • 全版本兼容:支持所有 Composer 版本(Composer 1 / 2 / 3),源码内部会根据检测到的 Airflow 版本自动切换不同的 gcloud 子命令。

该脚本是 composer/tools 目录下两个实用工具之一(另一个是环境迁移工具 composer_migrate.md),并已收录于 composer/tools/README.md 的入口索引中。如果你正在规划 Composer 2 → Composer 3 的迁移,迁移脚本内部同样复用了"逐个暂停/恢复 DAG"的策略,理解本脚本的实现思路有助于同时掌握两个工具的行为。

前置条件:运行前必须逐项确认

根据文档 Prerequisites,脚本运行依赖以下五项条件,缺一不可:

  1. 授权与权限

    • 运行前需先执行gcloud auth login完成身份认证;
    • 需要具备访问 Composer 环境的权限,roles/composer.environmentAndStorageObjectAdmin角色即可满足要求;
    • 同时需要拥有该环境底层 GKE 集群的访问权限(container.*),因为脚本需要通过 kubectl 与集群通信。
  2. kubectl 可访问 Composer 环境的 GKE 集群

    • 建议先用kubectl get pods验证连通性;若该命令本身失败,脚本必然失败;
    • 典型错误特征:kubectl 输出Unable to connect to the server: dial tcp [IP ADDRESS]:443: connect: connection timed out,说明 kubectl 无法连接到 GKE 集群;
    • 如果你使用的是私有 IP 环境(Private Cluster),需要按官方文档先行配置私有集群的访问方式(如 Cloud Shell 代理),确保 kubectl 能直连集群 API Server。
  3. 环境版本:支持所有 Composer 版本,脚本会自动适配。

  4. 本机工具链:Python 3.6 或更高版本、gcloud CLI、kubectl 三者缺一不可(源码依赖subprocess调用 gcloud 与 kubectl 相关命令,Python 内置的argparse、json、re、logging模块则负责参数解析与日志输出,见 composer_dags.py 导入区)。

  5. 环境健康:运行前请确认 Composer 环境的 "Environment health" 与 "Database health" 指标正常。若环境不健康(如数据库异常、调度器故障),应先修复环境再执行脚本,否则批量操作可能产生不可预期的半完成状态。

已知限制:使用前需了解的行为边界

文档在 Limitations 中明确了两条限制:

  1. gcloud 命令失败无错误处理机制:脚本当前没有针对 gcloud 命令执行失败的兜底错误处理。不过从源码看,底层_run_shell_command_locally_once在子进程返回非零退出码时会记录错误日志并直接sys.exit(1)终止(见 composer_dags.py#L80-L86),因此失败会"快速失败"而非静默继续。
  2. 操作后不做 DAG 校验:脚本不会在批量暂停/恢复之后验证 DAG 的实际状态是否全部符合预期,需要用户自行通过 Airflow UI 或gcloud composer environments run ... dags list复核。

此外,从源码实现可以补充两个值得注意的行为细节:

  • 单个 DAG 的 pause/unpause 若首次执行返回码为 1,脚本会自动重试一次(见 pause_dag 与 unpause_dag),两次都失败才放弃并打印 "Unable to pause/unpause DAG" 日志;
  • 文档声称"除airflow_monitoring外全部操作",但从 main 函数源码 看,pause 分支显式跳过airflow_monitoring,而 unpause 分支并未显式排除它,即恢复操作会覆盖包括airflow_monitoring在内的所有 DAG。实际使用时建议留意这一差异。

命令行参数详解

脚本通过argparse解析参数(见 parse_arguments),参数清单如下:

参数类型是否必填默认值说明
--operationstr是无操作类型,取值仅限pause或unpause(choices强约束,传其他值直接报错)
--projectstr是无项目名称(Project ID)
--environmentstr是无Composer 环境名称(Environment Name)
--locationstr是无环境所在区域(Region,如europe-west4)
--sdk_endpointstr否https://composer.googleapis.com/Composer SDK API 端点,用于覆盖 API 端点(如搭配 VPC-SC / 代理场景),通过环境变量CLOUDSDK_API_ENDPOINT_OVERRIDES_COMPOSER注入到每条 gcloud 命令前

其中--sdk_endpoint属于进阶参数,绝大多数场景无需修改;保持默认值即可命中公共 Composer API 端点。

使用方法:暂停 / 恢复全部 DAG

文档在 Usage 中给出了标准调用格式与两个示例。通用命令模板为:

python3 composer_dags.py \ --project [PROJECT NAME] \ --environment [SOURCE ENVIRONMENT NAME] \ --location [REGION] \ --operation [OPERATION]

暂停全部 DAG 示例:

python3 composer_dags.py \ --project my-project \ --environment my-airflow-1-composer-environment \ --location europe-west4 \ --operation pause

恢复全部 DAG 示例:

python3 composer_dags.py \ --project my-project \ --environment my-airflow-1-composer-environment \ --location europe-west4 \ --operation unpause

脚本入口位于 main,执行流程为:parse_arguments()解析参数 → 调用main(...)→ 按返回值exit。main内部的实际流水线分为四步(见 composer_dags.py#L155-L204):

  1. 调用describe_environment获取环境 JSON 描述,打印环境名称;
  2. 从 JSON 中的config.softwareConfig.imageVersion字段解析出 Airflow 版本号;
  3. 调用get_list_of_dags列出环境内全部 DAG,并打印List of dags : [...]日志;
  4. 根据--operation逐条遍历执行 pause/unpause。

运行时脚本会输出带时间戳的 DEBUG 级日志(日志格式在 logging.basicConfig 中定义),包括每条被执行的 shell 命令(Executing shell command: ...)、环境信息、镜像版本与 DAG 列表,便于跟踪执行进度和排查问题。

源码级原理:版本识别、DAG 列举与命令构造

1. 自动识别 Airflow 版本并切换子命令

脚本通过正则表达式COMPOSER_AF_VERSION_RE从环境镜像版本号(形如composer-2.x.x-airflow-2.x.x)中抽取 Airflow 版本:

COMPOSER_AF_VERSION_RE = re.compile( "composer-(\d+)(?:\.(\d+)\.(\d+))?.*?-airflow-(\d+)\.(\d+)\.(\d+)" )

(见 composer_dags.py#L35-L37)

在 main 中,environment_info["config"]["softwareConfig"]["imageVersion"]经正则匹配后取第 4、5、6 组捕获值,构造成(major, minor, patch)形式的airflow_version元组。

2. 新旧两套 gcloud 子命令的分流

脚本兼容 Airflow 1.x 与 2.x 的关键在于子命令命名差异,由airflow_version < (2, 0, 0)决定:

操作Airflow < 2.0(旧式子命令)Airflow >= 2.0(新式子命令)
列举 DAGlist_dagsdags list
暂停 DAGpausedags pause
恢复 DAGunpausedags unpause

对应实现位置:get_list_of_dags 的子命令选择、pause_dag 的子命令选择、unpause_dag 的子命令选择。命令统一通过gcloud composer environments run <environment> --project=<project> --location=<location> <sub_command>组装,并对单个 DAG 追加-- <dag_id>参数传入。

3. DAG 列表的解析策略

get_list_of_dags(见 composer_dags.py#L40-L66)对两种输出格式分别解析:

  • 旧版本(Airflow < 2.0):命令输出按空白符split(),定位DAGS标记后截取其后第 3 个元素到倒数第 2 个元素作为 DAG 名列表;
  • 新版本(Airflow >= 2.0):按行遍历输出,用^[a-zA-Z].*正则过滤出以字母开头的行,取每行第一个字段(DAG ID),并丢弃首行表头(list_of_dags[1:])。

4. 底层 shell 执行与失败处理

所有 gcloud 命令都由_run_shell_command_locally_once(见 composer_dags.py#L68-L87)统一执行:使用subprocess.Popen(..., shell=True)运行,通过communicate()收集输出;若返回码非零,则记录Failed to run shell command ...错误日志并sys.exit(1)终止脚本。

故障排查指南

文档在 Troubleshooting 中给出了三层排查路径:

  1. 逐项核对前置条件:确认已gcloud auth login授权、权限齐全、kubectl 可达、工具链完整、环境健康;
  2. 权限相关报错:如果错误疑似与权限有关,请重点复查是否同时具备 Composer 环境的访问权限(roles/composer.environmentAndStorageObjectAdmin)与 GKE 集群的container.*权限;
  3. 联系支持:若问题仍未解决,可向 Google Cloud 支持渠道求助,联系时务必附带脚本的完整输出日志(包括Executing shell command、镜像版本、DAG 列表以及具体的失败命令与输出),这能大幅缩短定位时间。

结合源码可以补充两条实用排查技巧:

  • 若日志中出现Failed to run shell command ... details: ...后脚本退出码为 1,说明某条 gcloud 命令执行失败,可直接复制日志中打印的完整命令到本地终端手动执行,以区分是权限问题、网络问题还是参数问题;
  • 若kubectl get pods已可连通,但脚本仍失败,优先检查环境健康状态——脚本在 pause/unpause 前会先执行describe与list_dags,这两个步骤的失败信息通常会直接暴露根因。

总结

composer_dags.py是一个轻量、自包含的运维工具,通过"描述环境 → 解析 Airflow 版本 → 列举 DAG → 逐条执行"四步流水线,将 Composer 环境整体的 DAG 暂停/恢复操作从数十次手工点击简化为一条命令行。其价值体现在:全 Composer 版本兼容(自动切换新旧子命令)、自动跳过airflow_monitoring(暂停时)、失败自动重试一次、全程 DEBUG 日志可审计。对于需要批量管理 DAG 状态的运维与迁移场景,可直接参照本文参数表与命令示例运行;如需进一步了解配套的 Composer 2 → 3 环境迁移工具,可阅读同目录下的 composer_migrate.md。

  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

相关推荐

上一篇:保姆级上手攻略:如何零网络玩转Zwift离线版?3步告别网络焦虑
下一篇:零网络畅骑Zwift:Zoffline离线版完整指南,4种安装方法与进阶玩法一次讲透

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/3 2:18:35

Tool Graph Retriever: Exploring Dependency Graph-based Tool Retrieval for Large Language Models

文章主要内容和创新点 主要内容 本文针对大型语言模型(LLM)驱动的AI代理在工具检索中存在的问题,提出了一种名为Tool Graph Retriever(TGR,工具图检索器) 的方法。 背景:随着AI代理配备的工具数量激增,受限于模型上下文长度,需通过检索筛选工具。现有方法主要依赖用…

作者头像 李华