Apache Airflow airflowctl 命令行与环境变量完全参考指南
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
导读
本文围绕 Apache Airflow 新版 CLI 工具airflowctl(位于仓库 airflow-ctl/ 目录)展开,系统讲解其命令行接口的完整结构、核心命令分组、参数体系,以及AIRFLOW_CLI_*系列环境变量的作用与源码级实现原理。读完本文,你将掌握如何用 airflowctl 完成认证登录、Dag 管理、任务诊断、连接/变量/资源池导入导出等日常运维操作,理解其"命令由 OpenAPI 操作自动生成 + 手写命令合并"的双轨架构,并能正确配置环境变量以控制 API 重试、调试模式与环境选择。
airflowctl 是什么:Airflow 的命令行操作平台
airflowctl 是 Apache Airflow 为 Airflow 3 提供的下一代命令行接口。与传统的本地airflowCLI 不同,airflowctl 面向通过 Airflow REST API 远程管理元数据库的场景:它可以对 Dag 执行暂停/恢复、查询下次调度时间、查看 Dag run 状态、诊断任务未调度的原因、检查调度器/触发器/处理器等 Job 存活状态,还支持从旧版 CLI 导出的文件导入连接、变量与资源池。从 airflow-ctl/src/airflowctl/ctl/cli_config.py 可以看到,它提供auth、config、connections、dags、jobs、pools、tasks、variables、version等命令分组,覆盖了"登录认证 → 运维管理 → 开发调试"的完整链路。
airflowctl 的命令行极其丰富,核心定位是"对 Dag 进行多种操作、启动服务、支持开发与测试"。执行入口在 airflow-ctl/src/airflowctl/main.py:先通过cli_parser.get_parser()构建参数解析器,再借助argcomplete提供 shell 自动补全,最后由safe_call_command统一执行命令并处理异常。
命令体系总览:分组命令与操作命令
airflowctl 的命令树由两类节点构成(定义见 airflow-ctl/src/airflowctl/ctl/cli_config.py):
- ActionCommand(操作命令):单个可执行命令,包含名称、帮助文本、回调函数与参数列表;
- GroupCommand(分组命令):包含一组子命令的容器,例如
dags下挂着clear、next-execution、pause、state、unpause等子命令。
所有命令统一注册在core_commands列表中,最终合并顺序为:由操作类自动生成的命令 + 手写核心命令 + 统一注入--api-token参数(见 airflow-ctl/src/airflowctl/ctl/cli_config.py 的merge_commands与add_auth_token_to_all_commands)。
命令分组速查表
| 分组 | 用途 | 子命令示例 |
|---|---|---|
auth | 认证与会话管理 | login、list-envs、token |
config | 配置查看与迁移检查 | lint |
connections | 连接导入 | import |
dags | Dag 运维 | clear、next-execution、pause、state、unpause |
jobs | Job 存活检查 | check |
pools | 资源池管理 | import、export |
tasks | 任务实例诊断 | failed-deps、states-for-dag-run |
variables | 变量导入 | import |
version | 版本信息 | 直接执行 |
此外,CommandFactory会根据 airflow-ctl/src/airflowctl/api/operations.py 中定义的*Operations类(如DagsOperations、TaskInstancesOperations、ConnectionsOperations、VariablesOperations、PoolsOperations等),通过 AST 解析自动生成一批"操作命令",再与手写命令合并。这意味着 airflowctl 的绝大部分能力直接对齐 Airflow 的 OpenAPI 接口,属于"API 操作的一等 CLI 客户端"。
通用参数
--output/-o:输出格式,取值table、json、yaml、plain,默认json(定义见 airflow-ctl/src/airflowctl/ctl/cli_config.py)。-e/--env:指定运行环境,默认production。--api-token:所有命令都会自动注入的认证令牌参数。--preview:由 airflow-ctl/src/airflowctl/ctl/cli_parser.py 为每个命令提供的帮助预览动作。
在非 TTY(管道)环境下,AirflowConsole会把输出宽度固定为 200 列以保证整段输出可被管道完整读取(见 airflow-ctl/src/airflowctl/ctl/console_formatting.py),非常适合脚本化调用。
认证与多环境管理(auth)
login:登录并持久化凭证
airflowctl auth login支持三种认证方式,逻辑见 airflow-ctl/src/airflowctl/ctl/commands/auth_command.py:
- 用户名 + 密码:直接通过
--username/--password传入,或由 CLI 交互式提示输入;调用/auth端点换取 JWT token 后保存; - JWT Token:通过
--api-token或环境变量AIRFLOW_CLI_TOKEN传入,保存 token; - 若未提供任何凭证且 stdin 不是终端,则报错并提示正确的用法后退出。
# 方式一:用户名密码登录(会提示输入,凭证保存到 keyring) airflowctl auth login --username admin --password secret --api-url http://localhost:8080 # 方式二:Token 登录(token 也可来自 AIRFLOW_CLI_TOKEN 环境变量) airflowctl auth login --api-token <JWT> --api-url http://localhost:8080 # 跳过系统 keyring,仅保存 API URL 配置 airflowctl auth login --api-token <JWT> --skip-keyring登录后,凭证按环境名维度保存:API URL 写入AIRFLOW_HOME(默认~/airflow)下的<env>.json配置文件,token 默认存入系统 keyring(--skip-keyring可跳过)。auth login使用的 base URL 是{api_url}/auth(见 airflow-ctl/src/airflowctl/api/client.py),与普通命令的/api/v2前缀区分。
list-envs:列出已登录环境
airflowctl auth list-envs扫描AIRFLOW_HOME下的*.json配置文件(跳过debug_creds_*与*_generated.json),输出每个环境名、API URL 与认证状态(authenticated / not authenticated / keyring unavailable 等),见 airflow-ctl/src/airflowctl/ctl/commands/auth_command.py。
token:生成并打印 JWT
airflowctl auth token使用用户名密码向认证端点换取 token 并打印到 stdout,适合在脚本中动态获取凭证:
airflowctl auth token --username admin --password secretDag 运维命令(dags)
dags分组是 airflowctl 的核心运维面,实现位于 airflow-ctl/src/airflowctl/ctl/commands/dag_command.py。
pause / unpause:暂停与恢复
通过PATCH /dags/{dag_id}更新is_paused标志(DAGPatchBody(is_paused=...)),成功后打印绿色提示并输出 Dag 详情:
airflowctl dags pause my_dag airflowctl dags unpause my_dag -o tablenext-execution:查询下次调度时间
读取 Dag 的next_dagrun_logical_date、next_dagrun_data_interval_start、next_dagrun_data_interval_end、next_dagrun_run_after四个字段并输出;若无即将到来的运行则提示No upcoming run scheduled:
airflowctl dags next-execution my_dagstate:查看 Dag run 状态
必须二选一传入run_id或--logical-date(同时传或都不传会报错)。--logical-date必须包含时区偏移,按精确匹配查询 Dag run,最终输出状态值(以及conf内容):
airflowctl dags state my_dag run_id=scheduled__2026-09-10T00:00:00+00:00 airflowctl dags state my_dag --logical-date 2026-09-10T00:00:00+00:00clear:清理任务实例
airflowctl dags clear支持三种互斥的 Dag run 选择方式,参数校验逻辑见 airflow-ctl/src/airflowctl/ctl/commands/dag_command.py:
--run-id:按运行 ID 精确选择;--partition-key:按分区键选择(Airflow 3 新增概念);--partition-date-start+--partition-date-end:按分区日期闭区间窗口选择,日期按 Dag 调度时区的本地日历日解释,忽略时间部分。
辅助参数:
| 参数 | 作用 |
|---|---|
-f/--only-failed | 只清理失败的任务实例 |
-r/--only-running | 只清理运行中的任务实例(与--only-failed互斥) |
-y/--yes | 跳过确认提示 |
执行时会先列出将要清理的 Dag run 并交互确认(Clear task instances for these Dag runs? [y/N]),确认后逐 run 调用POST /dags/{dag_id}/clearTaskInstances(dry_run=False),汇总被清理的实例数:
airflowctl dags clear my_dag --partition-date-start 2026-09-01 --partition-date-end 2026-09-07 -y任务诊断命令(tasks)
tasks分组提供两个面向调度的诊断命令(实现见 airflow-ctl/src/airflowctl/ctl/commands/task_command.py):
airflowctl tasks failed-deps DAG_ID TASK_ID [run_id] [--logical-date ...] [--map-index N]:返回任务实例未满足的依赖项,即从调度器视角解释"为什么这个任务实例没有被调度、排队并最终被执行"。--map-index用于指定动态任务映射的索引,默认-1。airflowctl tasks states-for-dag-run DAG_ID [run_id] [--logical-date ...]:输出某次 Dag run 内所有任务实例的状态。run 的选择同样遵循"run_id 与 --logical-date 二选一"的约束(相关 Arg 定义见 airflow-ctl/src/airflowctl/ctl/cli_config.py)。
Job 存活检查(jobs)
airflowctl jobs check用于确认调度类进程是否存活,参数见 airflow-ctl/src/airflowctl/ctl/cli_config.py:
| 参数 | 说明 |
|---|---|
--job-type | 过滤 Job 类型,可选SchedulerJob、TriggererJob、DagProcessorJob |
--hostname | 按主机名过滤 |
--local | 只显示本机 Job |
--limit | 检查最近 N 个 Job,默认 1,设为 0 表示不限 |
--allow-multiple | 即使找到多个匹配的存活 Job 也视为成功 |
典型用法:检查调度器是否存活。
airflowctl jobs check --job-type SchedulerJob --limit 3连接、变量与资源池的导入导出
旧版本地 CLI 导出的连接、变量、资源池文件可以无缝迁移到 airflowctl:
# 导入连接(从旧版 airflow connections export 导出的 JSON) airflowctl connections import connections.json # 导入变量 airflowctl variables import variables.json # 导入/导出资源池 airflowctl pools import pools.json airflowctl pools export pools_out.json三者都支持-a/--action-on-existing-key,用于指定实体已存在时的行为,取值overwrite(默认)、fail、skip(见 airflow-ctl/src/airflowctl/ctl/cli_config.py)。
config lint:Airflow 2 → 3 配置迁移检查
airflowctl config lint用于在从 Airflow 2 迁移到 Airflow 3 时对配置变更做静态检查,参数定义见 airflow-ctl/src/airflowctl/ctl/cli_config.py:
| 参数 | 说明 |
|---|---|
--section | 要检查的配置节 |
--option | 要检查的配置项 |
--ignore-section | 忽略的配置节 |
--ignore-option | 忽略的配置项 |
-v/--verbose | 输出详细结果,包括被忽略的节与项 |
version:版本信息
airflowctl version # 仅显示本地 airflowctl 版本 airflowctl version --remote # 额外拉取远端 Airflow 服务器版本环境变量参考(AIRFLOW_CLI_*)
airflowctl 通过一组AIRFLOW_CLI_*环境变量控制运行行为,均在 airflow-ctl/src/airflowctl/api/client.py 与 airflow-ctl/src/airflowctl/ctl/cli_config.py 中被读取。
AIRFLOW_CLI_TOKEN
用于与 Airflow API 认证的令牌,仅在未使用其他认证方式(如用户名密码)时需要。它同时是auth login的 token 来源(见 airflow-ctl/src/airflowctl/ctl/commands/auth_command.py),也会在get_client中作为兜底 token 注入(见 airflow-ctl/src/airflowctl/api/client.py)。优先级低于显式传入的--api-token参数。
export AIRFLOW_CLI_TOKEN="<your-jwt>" airflowctl dags pause my_dagAIRFLOW_CLI_ENVIRONMENT
指定 CLI 使用的环境名。在未设置时使用默认环境production。环境名会在 airflow-ctl/src/airflowctl/api/client.py 中被校验:包含/、\或..的环境名会被拒绝,以避免路径穿越,同时环境名直接决定了本地配置文件名(<env>.json)与 keyring 键名。
export AIRFLOW_CLI_ENVIRONMENT=staging airflowctl dags next-execution my_dag注意:-e/--env参数与--api-url会一起参与环境的选择;若未显式传入,AIRFLOW_CLI_ENVIRONMENT将决定使用哪份已保存的凭证与 API URL。
AIRFLOW_CLI_DEBUG_MODE
用于开启 CLI 调试模式。该模式会禁用 keyring 集成、改为把凭证保存到调试文件(AIRFLOW_HOME/debug_creds_<env>.json),并打印黄色警告提示凭证不安全(见 airflow-ctl/src/airflowctl/ctl/cli_config.py 与 airflow-ctl/src/airflowctl/api/client.py)。它仅适用于开发 airflowctl 或运行 API 集成测试的场景,官方明确提示:除非你清楚自己在做什么,否则不要使用。
# 调试模式下登录(凭证存入 debug_creds_<env>.json) AIRFLOW_CLI_DEBUG_MODE=true airflowctl auth login --api-token <JWT> # 完成后务必关闭 unset AIRFLOW_CLI_DEBUG_MODEAPI 重试三变量
当通过 Airflow API 调用失败时,以下三个变量控制重试策略(默认值及读取逻辑见 airflow-ctl/src/airflowctl/api/client.py):
| 变量 | 作用 | 默认值 |
|---|---|---|
AIRFLOW_CLI_API_RETRIES | API 调用失败后的重试次数 | 3 |
AIRFLOW_CLI_API_RETRY_WAIT_MIN | 两次重试之间的最小等待时间(秒) | 1 |
AIRFLOW_CLI_API_RETRY_WAIT_MAX | 两次重试之间的最大等待时间(秒) | 10 |
重试机制通过tenacity的@retry装饰器实现:仅对服务端 5xx 错误与网络请求异常(httpx.RequestError)进行重试,等待策略为随机指数退避(wait_random_exponential),见 airflow-ctl/src/airflowctl/api/client.py 与 airflow-ctl/src/airflowctl/api/client.py。
取值校验规则(源码级行为,见 airflow-ctl/src/airflowctl/api/client.py):
- 任意一个重试变量若不是整数,或超出范围(
AIRFLOW_CLI_API_RETRIES小于 1、两个等待变量为负数),会向stderr输出警告并回退到默认值; - 将变量设为空字符串等同于未设置,同样使用默认值(这也兼容了 Docker / Kubernetes 环境把"未配置"表达为空值的情况)。
# 调大重试次数并收紧等待窗口 export AIRFLOW_CLI_API_RETRIES=5 export AIRFLOW_CLI_API_RETRY_WAIT_MIN=2 export AIRFLOW_CLI_API_RETRY_WAIT_MAX=8命令解析与执行的底层原理
从 OpenAPI 操作到 CLI 命令的自动生成
airflowctl 的"操作命令"并非手写,而是由CommandFactory在运行时生成(见 airflow-ctl/src/airflowctl/ctl/cli_config.py):
- 用
ast解析 airflow-ctl/src/airflowctl/api/operations.py,提取所有*Operations类的公有方法签名(参数、必填项、默认值、返回类型); - 根据参数类型生成对应
Arg:必填的原始类型参数暴露为位置参数,可选参数与布尔参数暴露为--flag/--no-flag风格选项;布尔参数统一使用argparse.BooleanOptionalAction; - 非原始类型的 Pydantic datamodel 参数会被展开为其字段,例如
ClearTaskInstancesBody的字段逐一映射为--task-ids、--dry-run等选项; - 属于
list/get/create/delete/update/trigger/add/edit/set/clear前缀的操作自动附加--output与-e/--env参数; - 最后通过
partial绑定操作信息生成回调,执行时动态实例化操作类并调用对应方法。
其中针对ClearTaskInstancesBody还有特殊处理:--task-ids既接受逗号分隔的 ID 列表,也接受 JSON 数组字面量(见 airflow-ctl/src/airflowctl/ctl/cli_config.py)。
命令分发与异常处理
用户敲下命令后,执行链路为:
- airflow-ctl/src/airflowctl/main.py 中
get_parser()构建 argparse 解析器(含argcomplete补全); - 解析出的
Namespace携带func回调(每个ActionCommand通过set_defaults(func=...)绑定,见 airflow-ctl/src/airflowctl/ctl/cli_parser.py); safe_call_command统一执行并分类处理异常(见 airflow-ctl/src/airflowctl/ctl/cli_config.py):凭证/连接/keyring/404 类异常直接报错退出;httpx协议错误与超时会给出"请检查服务器是否运行、API URL 是否正确"等针对性提示;服务端响应错误则提示结合--help检查命令与参数。
输出渲染
所有结果最终经AirflowConsole.print_as渲染(见 airflow-ctl/src/airflowctl/ctl/console_formatting.py),支持 json / yaml / table / plain 四种格式:JSON 与 YAML 用语法高亮输出,table 用 rich Table,plain 用tabulate生成可直接被管道消费的纯文本表格。非 TTY 场景自动放宽列宽,保证脚本解析稳定。
快速上手:完整运维流程示例
# 1. 配置 API 地址与登录(默认 http://localhost:8080) airflowctl auth login --username admin --password secret # 2. 查看当前已登录的环境与状态 airflowctl auth list-envs # 3. 检查调度器是否存活 airflowctl jobs check --job-type SchedulerJob # 4. 查看某 Dag 下次调度时间与当前 run 状态 airflowctl dags next-execution my_dag airflowctl dags state my_dag --logical-date 2026-09-10T00:00:00+00:00 # 5. 诊断某任务为何未被调度 airflowctl tasks failed-deps my_dag my_task --logical-date 2026-09-10T00:00:00+00:00 # 6. 清理失败任务并重启 Dag run airflowctl dags clear my_dag --partition-date-start 2026-09-01 --partition-date-end 2026-09-07 -y # 7. 批量导入连接与变量,完成迁移 airflowctl connections import connections.json --action-on-existing-key skip airflowctl variables import variables.json参考资料
- 官方 CLI 与环境变量参考:本文源头文档 airflow-ctl/docs/cli-and-env-variables-ref.rst
- 命令定义与参数体系:airflow-ctl/src/airflowctl/ctl/cli_config.py
- 参数解析器:airflow-ctl/src/airflowctl/ctl/cli_parser.py
- API 客户端与重试/环境变量实现:airflow-ctl/src/airflowctl/api/client.py
- 认证命令实现:airflow-ctl/src/airflowctl/ctl/commands/auth_command.py
- Dag 命令实现:airflow-ctl/src/airflowctl/ctl/commands/dag_command.py
- 输出渲染:airflow-ctl/src/airflowctl/ctl/console_formatting.py
- 执行入口:airflow-ctl/src/airflowctl/main.py
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考