- 示例工程
【免费下载链接】python-docs-samples
Code samples used on cloud.google.com
导读
本文围绕 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,脚本运行依赖以下五项条件,缺一不可:
授权与权限
- 运行前需先执行
gcloud auth login完成身份认证; - 需要具备访问 Composer 环境的权限,
roles/composer.environmentAndStorageObjectAdmin角色即可满足要求; - 同时需要拥有该环境底层 GKE 集群的访问权限(
container.*),因为脚本需要通过 kubectl 与集群通信。
- 运行前需先执行
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。
- 建议先用
环境版本:支持所有 Composer 版本,脚本会自动适配。
本机工具链:Python 3.6 或更高版本、gcloud CLI、kubectl 三者缺一不可(源码依赖
subprocess调用 gcloud 与 kubectl 相关命令,Python 内置的argparse、json、re、logging模块则负责参数解析与日志输出,见 composer_dags.py 导入区)。环境健康:运行前请确认 Composer 环境的 "Environment health" 与 "Database health" 指标正常。若环境不健康(如数据库异常、调度器故障),应先修复环境再执行脚本,否则批量操作可能产生不可预期的半完成状态。
已知限制:使用前需了解的行为边界
文档在 Limitations 中明确了两条限制:
- gcloud 命令失败无错误处理机制:脚本当前没有针对 gcloud 命令执行失败的兜底错误处理。不过从源码看,底层
_run_shell_command_locally_once在子进程返回非零退出码时会记录错误日志并直接sys.exit(1)终止(见 composer_dags.py#L80-L86),因此失败会"快速失败"而非静默继续。 - 操作后不做 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),参数清单如下:
| 参数 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
--operation | str | 是 | 无 | 操作类型,取值仅限pause或unpause(choices强约束,传其他值直接报错) |
--project | str | 是 | 无 | 项目名称(Project ID) |
--environment | str | 是 | 无 | Composer 环境名称(Environment Name) |
--location | str | 是 | 无 | 环境所在区域(Region,如europe-west4) |
--sdk_endpoint | str | 否 | 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):
- 调用
describe_environment获取环境 JSON 描述,打印环境名称; - 从 JSON 中的
config.softwareConfig.imageVersion字段解析出 Airflow 版本号; - 调用
get_list_of_dags列出环境内全部 DAG,并打印List of dags : [...]日志; - 根据
--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(新式子命令) |
|---|---|---|
| 列举 DAG | list_dags | dags list |
| 暂停 DAG | pause | dags pause |
| 恢复 DAG | unpause | dags 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 中给出了三层排查路径:
- 逐项核对前置条件:确认已
gcloud auth login授权、权限齐全、kubectl 可达、工具链完整、环境健康; - 权限相关报错:如果错误疑似与权限有关,请重点复查是否同时具备 Composer 环境的访问权限(
roles/composer.environmentAndStorageObjectAdmin)与 GKE 集群的container.*权限; - 联系支持:若问题仍未解决,可向 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
相关推荐
如何快速掌握GCViewer:全面解读Java GC暂停、Full GC与安全点暂停分析指南
如何快速掌握GCViewer:全面解读Java GC暂停、Full GC与安全点暂停分析指南 GCViewer是一款强大的Java垃圾回收日志分析工具,专为解析
开发工具性能剖析数据可视化x64dbg 脚本命令 pause 详解:暂停脚本执行、恢复机制与脚本状态机原理
x64dbg 脚本命令 pause 详解:暂停脚本执行、恢复机制与脚本状态机原理 本篇文章讲解 x64dbg 内置脚本引擎(simplescript)中的 pa
逆向工程调试器开发工具应用安全co源码剖析:理解Generator函数的暂停与恢复机制
co源码剖析:理解Generator函数的暂停与恢复机制 在JavaScript异步编程的世界中,co库是一个革命性的工具,它让Generator函数的威力得到
后端开发工具
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考