资讯动态

Apache Airflow CLI 实现指南:airflowctl 与 airflow 双 CLI 架构下的新命令开发与重连实践

发布时间:2026/9/11 23:33:36 来源:尧图企业网站定制
Apache Airflow CLI 实现指南airflowctl 与 airflow 双 CLI 架构下的新命令开发与重连实践【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读本文基于 Apache Airflow 官方贡献文档《CLI Implementation Guide》编写系统阐述 Airflow 项目在 AIP-94 指导下确立的双 CLI 架构随核心发行版捆绑的airflowCLI 与独立发行、完全通过 Public API 通信的airflowctlCLI。阅读本文后你将掌握三类 CLI 开发场景的完整落地路径——新增可通过 Public API 实现的命令、将存量airflow远程命令重连到 API、以及新增无 API 对应的 admin/local 命令——并理解底层 HTTP 客户端、操作层与命令行配置层的源码实现。背景Airflow 的双 CLI 架构Apache Airflow 目前同时发行两个命令行工具它们的分工在 27_cli_implementation_guide.rst 中被明确界定airflowairflow-core随核心发行版捆绑。同时承载两类命令——正在被内部重连rewire的存量远程命令以及没有 Public API 对应的 admin/local 命令如数据库管理、进程管理等。airflowctlairflow-ctl独立发行、独立打包的 CLI唯一的通信通道是运行中 Airflow 实例的 PublicCoreAPI绝不直接访问元数据库。这一架构建立在 AIP-81 引入的local与remote命令区分之上并由AIP-94进一步收口。AIP-94 为所有 CLI 工作确立了不可违反的两条规则新增命令只要能够通过 Public API 实现就必须只加在airflowctl中。同时把同一命令再添加到airflowCLI 是被劝阻的——那只会复制维护面而对用户毫无收益。存量airflowCLI 远程命令保留原位用户继续运行airflow dags list、airflow pools get等但其内部实现被重连为通过airflowctl的 HTTP 客户端调用 Public API而不是直接访问元数据库。两条规则共同达成的目标强制统一走 RBAC基于角色的访问控制、消除远程操作对元数据库的直接暴露、并消灭重复的代码路径。注AIP-94 的进展在 Apache 的 GitHub Projects #570、#571 中跟踪相关设计与社区讨论可参见 AIP-94 的 Confluence 页面本文聚焦仓库内可直接核验的代码实现外部跟踪链接不再展开。决策表命令该放到哪里原文档给出了一张决定“命令归属”的决策表这是所有 CLI 开发工作的起点务必先对照再动手场景放哪里备注新命令可通过 Public API 实现仅airflowctl除非核心内确有强烈需求否则不要加到airflowCLI新命令无法通过 Public API 实现airflowCLIadmin/local只能是 admin/local 命令见下文约束存量airflowCLI 命令可通过 Public API 实现airflowCLI → 委托给airflowctlHTTP 客户端重连禁止直接访问数据库禁止使用 SQLAlchemy/session存量airflowCLI 命令无法通过 Public API 实现airflowCLI保持不变保持纯airflow-core实现其中“无法通过 Public API 实现”的精确定义是该操作在 API 层面没有对应表示且本质上属于 admin/local 范畴——典型包括数据库 shell、schema 迁移、进程管理、需要直接基础设施访问的部署配置。新增命令可通过 Public API 实现只写 airflowctl这是最常规的开发路径。核心约束一句话只改airflow-ctl不碰airflow-core的数据模型与元数据库。涉及源码位置命令定义与实现airflow-ctl/src/airflowctl/ctl/commands/按命令组拆分模块如dag_command.py、pool_command.py、connection_command.py、variable_command.py等HTTP 客户端与操作层airflow-ctl/src/airflowctl/api/client.py 与 airflow-ctl/src/airflowctl/api/operations.py测试目录airflow-ctl/tests/实施步骤在airflow-ctl/src/airflowctl/ctl/commands/下找到合适的命令组模块在其中添加命令函数。通过airflowctl.api.clientHTTP 客户端与airflowctl.api.operations操作层调用 Public API不要importairflow-core的任何模型也不要触碰元数据库。如果所需的 API 端点尚不存在先补端点见 添加 API 端点指南。在airflow-ctl/tests/下添加测试。运行集成测试验证breeze testing airflow-ctl-integration-tests源码级印证一个命令的完整生命周期以 pools 命令为例观察 pool_command.py 的实现模式from airflowctl.api.client import NEW_API_CLIENT, Client, ClientKind, provide_api_client from airflowctl.api.datamodels.generated import PoolBody from airflowctl.ctl.console_formatting import AirflowConsole provide_api_client(kindClientKind.CLI) def export(args, api_client: Client NEW_API_CLIENT) - None: pools_response api_client.pools.list() pools_list [ { name: pool.name, slots: pool.slots, description: pool.description, include_deferred: pool.include_deferred, } for pool in pools_response.pools ] AirflowConsole().print_as(datapools_list, outputargs.output)关键点拆解provide_api_client(kindClientKind.CLI)装饰器这是所有airflowctl命令的统一入口。装饰器内部见 client.py在调用命令函数前自动完成凭证加载与客户端构建从AIRFLOW_HOME下的环境配置文件如production.json读取api_url从系统 keyring 读取 token也支持AIRFLOW_CLI_TOKEN环境变量或--api-token参数覆盖随后构造Client注入函数调用结束后自动关闭连接。文档明确注释该客户端仅允许在 mock 与测试场景下手动传入。api_client.pools.list()操作层对象。Client类通过lru_cache属性缓存暴露了dags、dag_runs、pools、variables、connections、tasks、task_instances、jobs、configs、assets、backfills、providers、version、xcom、plugins等 15 个操作对象见 client.py每个操作对象的方法直接对应 REST 端点。AirflowConsole().print_as(...)统一的控制台输出层支持table / json / yaml / plain四种格式--output/-o参数默认json定义于 cli_config.py。HTTP 客户端底层设计client.py 中的Client类继承自httpx.Client为 CLI 做了多项针对性增强Bearer 认证BearerAuth在请求头注入Authorization: Bearer tokenClientKind.NO_AUTH跳过认证。Base URL 自动拼接ClientKind.AUTH时指向{base_url}/auth其余指向{base_url}/api/v2。关联 ID每个请求通过add_correlation_id注入correlation-id头基于 uuid7便于分布式追踪。自动重试基于tenacity对 5xx 与httpx.RequestError指数退避重试。重试次数与等待区间可由环境变量调优AIRFLOW_CLI_API_RETRIES默认 3、AIRFLOW_CLI_API_RETRY_WAIT_MIN默认 1 秒、AIRFLOW_CLI_API_RETRY_WAIT_MAX默认 10 秒。错误标准化响应事件钩子raise_on_4xx_5xx会把 4xx/5xx 响应转换为ServerResponseError抛出同时保留原始 JSON 错误体方便定位问题。操作层设计operations.py 中的BaseOperations是所有操作类的基类其__init_subclass__会为所有可调用方法自动包裹_check_flag_and_exit_if_server_response_error装饰器默认exit_in_errorTrue一旦服务端返回错误立即抛出若遇httpx.ConnectError且为 Connection refused会给出友好提示 Connection refused. Is the API server running?。值得注意的通用方法execute_list它支持自动分页拉取——先请求limit条若total_entries超过当前页则循环按offset步进抓取直到取满并把分页字段拼接为完整的集合响应。DagRunOperations.list则展示了参数序列化细节datetime/date会被自动转为 ISO 格式字符串_serialize_query_param未指定的筛选条件自动剔除_build_query_params未传dag_id时默认使用~表示“全部 DAG”。命令行定义声明式配置 自动生成airflowctl的命令行参数体系是声明式的集中定义于 cli_config.pyArg类封装单个参数的 flags、help、action、default、choices、type 等元信息如ARG_OUTPUT、ARG_DAG_ID、ARG_ACTION_ON_EXISTING_KEYchoices 为overwrite / fail / skip默认overwrite。CommandFactory更为关键——它通过ast静态解析operations.py中每个操作方法的签名与类型注解自动生成对应的命令行参数_inspect_operations。必填的基本类型参数自动成为位置参数布尔参数自动生成为--flag/--no-flag风格Pydantic 数据模型参数则展开为其字段级别的选项list/get/create/delete/update/trigger/add/edit/set/clear等输出型操作自动追加--output与-e/--env。cli_parser.py将上述声明转换为真正的argparse解析器并按“Groups / Commands”分组展示 help支持--preview预览 action。换言之大多数新增命令甚至不需要手工编写参数解析逻辑——只要在operations.py中按规范添加一个带类型注解的操作方法CommandFactory就会为其生成 CLI 参数与调用绑定_create_func_map_from_operation负责把 Namespace 映射回操作方法调用。这从源码结构上保证了“CLI 与 API 一一对应、不产生漂移”的设计意图。重连存量 airflow CLI 命令改内部实现不动用户界面适用场景某个存量airflowCLI远程命令仍直连数据库需要改为经由 Public API 调用。硬性要求用户可见的命令名与参数完全不变。涉及源码位置airflow-core/src/airflow/cli/实施步骤将命令中的直接数据库访问替换为通过airflowctlHTTP 客户端airflowctl.api.client/airflowctl.api.operations的调用从命令中移除 SQLAlchemy 模型 import 与基于session的辅助函数。若所需 API 端点不存在先添加端点见 添加 API 端点指南。更新 airflow-core/tests/cli/ 下的测试改为 mock 或实际驱动 HTTP 客户端而非直接操作数据库。源码级印证pools 命令的重连现状以 airflow-core/src/airflow/cli/commands/pool_command.py 为例可以看到重连后的典型形态from airflow.api.client import get_current_api_client from airflow.cli.utils import deprecated_for_airflowctl deprecated_for_airflowctl(airflowctl pools list) suppress_logs_and_warning providers_configuration_loaded def pool_list(args): Display info of all the pools. api_client get_current_api_client() pools api_client.get_pools() _show_pools(poolspools, outputargs.output)三点关键证据get_current_api_client()命令不再 import 任何 SQLAlchemy 模型而是通过 Airflow 的 API 客户端抽象获取远程数据。deprecated_for_airflowctl(airflowctl pools list)重连的同时打上“迁移到 airflowctl”的弃用标记引导用户平滑过渡到新 CLI。在仓库中搜索airflowctl相关引用可发现provider_command.py、backfill_command.py、asset_command.py、jobs_command.py、connection_command.py、dag_command.py、task_command.py、config_command.py、variable_command.py、pool_command.py以及 cli_config.py 均已接入该模式——即重连工作已在核心命令组中铺开。此外legacy_commands.py 中维护了COMMAND_MAP如dags backfill: backfill create、webserver: api-server用于在用户输入已移除的旧命令时给出“请改用新命令”的提示保证迁移体验友好。新增 admin/local 命令无 Public API 对应的唯一去处仅当操作无法合理通过 Public API 暴露时才走此路径典型场景包括数据库 shell、schema 迁移、进程管理、部署期配置。涉及源码位置airflow-core/src/airflow/cli/实施步骤将命令添加到合适的 admin 命令组如db、config组对应 airflow-core/src/airflow/cli/commands/db_command.py、config_command.py等。在help字符串中标注(admin only)让用户明确知晓该命令需要直接的基础设施访问权限。在 airflow-core/tests/cli/ 下添加测试。这类命令保持纯airflow-core实现不涉及 HTTP 客户端也因此天然不受 AIP-94 第一条规则的约束——因为它们在 API 层根本没有对应端点。测试策略airflowctl 新命令单元/集成测试放在 airflow-ctl/tests/端到端集成测试可进一步参考 Airflow Ctl 测试文档。官方推荐的集成测试命令breeze testing airflow-ctl-integration-testsairflow 存量命令重连更新 airflow-core/tests/cli/ 下的既有测试将数据库依赖替换为对 HTTP 客户端的 mock 或真实调用。这与provide_api_client装饰器“仅测试场景可手动注入 api_client”的设计相呼应——测试中可显式传入 mock client 而无需真实网络。配套参考API 端点的添加CLI 重连与新增都经常依赖新端点端点的添加规范在 16_adding_api_endpoints.rst 中有完整说明核心要点端点位于api_fastapi/core_api/routes分public标准化、向后兼容、属于 Public API与ui仅供前端、不保证向后兼容两类优先做成 public 端点。用 FastAPI 路由 Pydantic 返回类型注解声明端点例如dags_router.get(/dags)并返回DagCollectionResponse。新 Pydantic 模型定义后只要被某个端点实际使用就会自动进入 OpenAPI spec持久化在v2-rest-api-generated.yaml运行prek --all-files即可让 prek hooks 更新生成文件。小结Apache Airflow 的 CLI 开发方向已经非常清晰以 Public API 为唯一边界把远程操作逐步收敛到airflowctl。新增功能一律先评估 API 可实现性——可实现则只写airflowctl且大多可借助CommandFactory自动生成参数存量远程命令则在不改变用户体验的前提下重连到 HTTP 客户端只有数据库迁移、进程管理等 admin/local 操作才继续留在airflowCLI。这一演进在代码层面由airflowctl.api.client/airflowctl.api.operations/airflowctl.ctl.cli_config三层结构强力支撑最终实现 RBAC 统一、数据库零直接暴露、代码路径单一化的长期目标。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价