
Apache Airflowresult装饰器让 DAG 返回结构化结果并通过 wait API 获取【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读Apache Airflow 在编排领域一直是调度与监控的代名词但对于调用方而言一个 DAG 运行结束后产出了什么往往是黑盒。本特性引入result装饰器将 TaskFlow 任务显式标记为 DAG 的结果任务result task并配合实验性的/dags/{dag_id}/dagRuns/{dag_run_id}/wait接口让调用方可以在 DagRun 结束后直接拿到任务返回值。读完本文你将掌握如何用result声明结果任务、理解dag装饰函数返回XComArg时的自动标记机制以及结果数据从任务执行、XCom 落库到 API 响应的完整链路。特性概述本特性对应 newsfragment 64563.feature.rst包含三个相互关联的能力新增result装饰器用于把 TaskFlow 任务标记为 DAG 的结果任务使用dag时从装饰函数中直接返回某个任务的XComArg也会自动把该任务标记为结果任务结果任务的返回值默认包含在GET /dags/{dag_id}/dagRuns/{dag_run_id}/wait的响应中——当result查询参数未被显式设置时。从迁移文件 0110_3_3_0_xcom_dag_result.py 的命名可以推断该特性随 Airflow 3.3.0 引入核心是给 XCom 模型新增dag_result布尔列用于在数据库层面标记这条 XCom 属于 DAG 结果。用result标记结果任务result必须叠加在task之上使用语法如下from airflow.sdk import result, task result task def emit_values(): something ... return something其实现位于 task-sdk/src/airflow/sdk/definitions/decorators/init.pydef result(t: C) - C: Mark a task as returning the dags result. This must be used *on top of* a task decorator like this:: result task def emit_values(): ... if not is_decorated_task(t): raise TypeError(result must be used on top of a task-decorated function) t.returns_dag_result True return t关键细节如果result没有用在被task装饰过的函数上会立即抛出TypeError提示必须叠加在task装饰器之上其本质是给装饰后的任务对象设置returns_dag_result True任务对象上该字段的默认值为False见 task-sdk/src/airflow/sdk/bases/decorator.py装饰器会原样返回任务对象因此可以继续参与依赖编排或.expand()映射展开标记行为在派生/覆盖属性时也会保留——对应测试 task-sdk/tests/task_sdk/definitions/decorators/test_result.py 验证了returns_dag_result被正确标记且在属性覆盖后依然保留。dag中返回XComArg自动标记除了显式使用result当 DAG 用dag装饰器定义时被装饰函数若返回某个任务的XComArg该任务会被自动设为结果任务from airflow.sdk import dag, task dag(scheduleNone, start_date...) def my_dag(): task def generate(): return {key: value} return generate() my_dag()这一逻辑由 DAG 类上的add_result方法驱动见 task-sdk/src/airflow/sdk/definitions/dag.pydef add_result(self, xcom_arg: X) - X: if not _is_valid_dag_result(xcom_arg): raise ValueError(Only plain return value can be used as dag result) xcom_arg.operator.returns_dag_result True return xcom_arg合法性校验函数_is_valid_dag_resulttask-sdk/src/airflow/sdk/definitions/dag.py要求该XComArg必须是普通非 mapped 派生返回值且 key 必须是 XCom 的XCOM_RETURN_KEY否则抛出ValueErrordef _is_valid_dag_result(value: Any) - TypeIs[PlainXComArg]: from airflow.sdk.bases.xcom import BaseXCom from airflow.sdk.definitions.xcom_arg import PlainXComArg return isinstance(value, PlainXComArg) and value.key BaseXCom.XCOM_RETURN_KEYdag装饰器在执行完被装饰函数后会检查返回值task-sdk/src/airflow/sdk/definitions/dag.pyr f(**f_kwargs) if _is_valid_dag_result(r): log.debug( Automatically adding function return value %r as result for dag %s, r, dag_obj.dag_id, ) dag_obj.add_result(r)对应的单元测试清晰地划定了行为边界task-sdk/tests/task_sdk/definitions/test_dag.pytest_ignore_function_resultdag函数返回普通值如return 123时不会触发add_result任务保持returns_dag_result is Falsetest_function_result_set_to_xcom_argreturn return_num(123)时returns_dag_result变为True。底层链路从任务执行到结果落库结果任务的标记最终体现在任务执行时的 XCom 推送环节。Task SDK 执行器在推送返回值 XCom 时会把任务的returns_dag_result一并写入见 task-sdk/src/airflow/sdk/execution_time/task_runner.pydef _xcom_push(ti, key, value, *, mapped_lengthNone): XCom.set( keykey, valuevalue, dag_idti.dag_id, task_idti.task_id, run_idti.run_id, map_indexti.map_index, dag_resultti.task.returns_dag_result, _mapped_lengthmapped_length, )对应到数据库层面XCom 模型新增了dag_result列airflow-core/src/airflow/models/xcom.pydag_result: Mapped[bool | None] mapped_column(Boolean, nullableTrue, defaultFalse)迁移脚本 0110_3_3_0_xcom_dag_result.py 通过batch_op.add_column(sa.Column(dag_result, sa.Boolean, nullableTrue))完成加列并提供了对应的降级脚本。在 Task SDK 执行 API 一侧POST /xcoms也新增了dag_result布尔查询参数见 airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py与任务运行时的推送保持一致。通过 wait API 获取结果接口形态GET /dags/{dag_id}/dagRuns/{dag_run_id}/wait是一个实验性端点在 OpenAPI 规范中被标记为experimental说明可能在没有预告的情况下变更或移除见 v2-rest-api-generated.yaml。其参数如下参数位置必填说明dag_idpath是DAG 标识dag_run_idpath是DagRun 标识intervalquery是轮询 DagRun 状态的间隔秒数必须大于 0exclusiveMinimum: 0.0resultquery否指定要收集结果 XCom 的任务 id可重复设置未设置时默认返回 DAG 中声明的结果任务即result或dag返回的 XComArg返回值流式 NDJSON 响应成功响应以换行分隔的 JSONNDJSON流式返回每一行是一个 JSON 对象代表 DagRun 的当前状态。OpenAPI 中的响应示例{state: running} {state: success, results: {op: 42}}即运行未结束时只返回state运行结束后附带results字段键为任务 id值为该任务返回值 XCom 的内容。服务端实现核心逻辑位于DagRunWaiter类airflow-core/src/airflow/api_fastapi/core_api/services/public/dag_run.py其wait方法按interval秒轮询 DagRun 状态并持续产出 NDJSON 行async def wait(self) - AsyncGenerator[str, None]: yield await self._serialize_response(dag_run : await self._get_dag_run()) yield \n while dag_run.state not in State.finished_dr_states: await asyncio.sleep(self.interval) yield await self._serialize_response(dag_run : await self._get_dag_run()) yield \n结果收集的关键在_serialize_xcoms当result_task_ids is None调用方未显式传result时查询该 DagRun 全部XCOM_RETURN_KEY且dag_result.is_(True)的 XCom即返回 DAG 作者声明的结果任务当调用方显式传了result时按指定的task_ids精确过滤结果统一按task_id, map_index排序以保证 mapped 任务的结果顺序稳定执行顺序本身不保证非 mapped 任务若只有一条 XCom 则解包为单个值mapped 任务则聚合成按map_index排序的列表。路由处理器wait_dag_run_until_finishedairflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py负责参数解析与权限校验若 DagRun 不存在返回 404若用户无 XCom 读取权限且未显式请求结果时会静默降级为不返回任何 XCom 结果result_task_ids []若显式请求了result但无权限则返回 403。三种调用方式对照场景请求行为不传resultDAG 声明了结果任务GET /wait?interval1默认返回result/dag返回标记的任务结果不传resultDAG 未声明结果任务GET /wait?interval1仅返回state无results字段显式指定任务GET /wait?interval1resulttask_1resulttask_2按指定任务收集可覆盖 DAG 作者声明映射任务的聚合行为result与动态任务映射mapping结合时所有映射实例的返回值会被聚合成一个按map_index排序的列表。单元测试 test_dag_run.py 完整验证了该行为def test_collect_mapped_task_dag_result(self, test_client, dag_maker, session): XComs from a mapped result task are aggregated into a list ordered by map_index. with dag_maker(dag_mapped_result): result task(task_ida) def double(v): return v * 2 mapped double.expand(v[1, 2]) ... assert response.json() {state: DagRunState.SUCCESS, results: {a: [2, 4]}}测试表明单个结果任务a映射展开后results中a的值为[2, 4]按 map_index 顺序即1*2与2*2。这与_serialize_xcoms中_group_xcoms的分组逻辑一致mapped 任务map_index 0的所有 XCom 值以列表返回非 mapped 任务解包为单值。权限、错误与边界TestWaitDagRun测试类test_dag_run.py系统性地覆盖了接口的各类边界401未认证客户端直接返回 401403无 DAG 访问权限返回 403有 RUN 权限但无 XCOM 权限时显式请求result返回 403未显式请求则降级为不返回结果状态码仍为 200404DagRun 不存在返回 404422缺少必填的interval参数返回 422隐式返回值DAG 声明了结果任务时{state: ..., results: {task_2: result_2}}DAG 未声明结果任务时仅返回{state: ...}显式返回值?resulttask_1可收集非结果任务?resulttask_2可收集结果任务均返回对应任务的返回值 XCom。此外权限校验采用了双重授权设计路由依赖先校验 RUN 访问权限处理器内再以相同的 team 解析方式校验 XCOM 访问权限避免不同粒度的校验对 team-aware 认证管理器产生不一致的授权判断。小结result装饰器与 wait API 构成了 Airflow 3.3.0 中结果感知的 DAG 执行闭环声明侧result叠加在task之上显式声明或通过dag函数返回XComArg隐式声明统一落到操作符的returns_dag_result标志执行侧Task SDK 推送返回值 XCom 时写入dag_resultTrue模型与迁移在数据库层面持久化该标记消费侧实验性 wait 端点以 NDJSON 流式返回 DagRun 状态与结果未显式指定result参数时默认返回 DAG 作者声明的结果任务mapped 结果按 map_index 聚合为列表。对于需要以编程方式触发并等待 Airflow DAG 完成的调用方如 CI/CD、数据平台上层编排这提供了一种无需轮询 XCom 明细即可获取 DAG最终产出的标准化方式。想深入验证行为可阅读 test_dag.py、test_result.py 与 test_dag_run.py 中的对应测试。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考