構(gòu)化結(jié)果并通過 wait API 獲取)
Apache Airflowresult裝飾器讓 DAG 返回結(jié)構(gòu)化結(jié)果并通過 wait API 獲取【免費下載鏈接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows項目地址: https://gitcode.com/GitHub_Trending/ai/airflow導讀Apache Airflow 在編排領(lǐng)域一直是調(diào)度與監(jiān)控的代名詞但對于調(diào)用方而言一個 DAG 運行結(jié)束后產(chǎn)出了什么往往是黑盒。本特性引入result裝飾器將 TaskFlow 任務顯式標記為 DAG 的結(jié)果任務result task并配合實驗性的/dags/{dag_id}/dagRuns/{dag_run_id}/wait接口讓調(diào)用方可以在 DagRun 結(jié)束后直接拿到任務返回值。讀完本文你將掌握如何用result聲明結(jié)果任務、理解dag裝飾函數(shù)返回XComArg時的自動標記機制以及結(jié)果數(shù)據(jù)從任務執(zhí)行、XCom 落庫到 API 響應的完整鏈路。特性概述本特性對應 newsfragment 64563.feature.rst包含三個相互關(guān)聯(lián)的能力新增result裝飾器用于把 TaskFlow 任務標記為 DAG 的結(jié)果任務使用dag時從裝飾函數(shù)中直接返回某個任務的XComArg也會自動把該任務標記為結(jié)果任務結(jié)果任務的返回值默認包含在GET /dags/{dag_id}/dagRuns/{dag_run_id}/wait的響應中——當result查詢參數(shù)未被顯式設(shè)置時。從遷移文件 0110_3_3_0_xcom_dag_result.py 的命名可以推斷該特性隨 Airflow 3.3.0 引入核心是給 XCom 模型新增dag_result布爾列用于在數(shù)據(jù)庫層面標記這條 XCom 屬于 DAG 結(jié)果。用result標記結(jié)果任務result必須疊加在task之上使用語法如下from airflow.sdk import result, task result task def emit_values(): something ... return something其實現(xiàn)位于 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關(guān)鍵細節(jié)如果result沒有用在被task裝飾過的函數(shù)上會立即拋出TypeError提示必須疊加在task裝飾器之上其本質(zhì)是給裝飾后的任務對象設(shè)置returns_dag_result True任務對象上該字段的默認值為False見 task-sdk/src/airflow/sdk/bases/decorator.py裝飾器會原樣返回任務對象因此可以繼續(xù)參與依賴編排或.expand()映射展開標記行為在派生/覆蓋屬性時也會保留——對應測試 task-sdk/tests/task_sdk/definitions/decorators/test_result.py 驗證了returns_dag_result被正確標記且在屬性覆蓋后依然保留。dag中返回XComArg自動標記除了顯式使用result當 DAG 用dag裝飾器定義時被裝飾函數(shù)若返回某個任務的XComArg該任務會被自動設(shè)為結(jié)果任務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方法驅(qū)動見 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合法性校驗函數(shù)_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裝飾器在執(zhí)行完被裝飾函數(shù)后會檢查返回值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函數(shù)返回普通值如return 123時不會觸發(fā)add_result任務保持returns_dag_result is Falsetest_function_result_set_to_xcom_argreturn return_num(123)時returns_dag_result變?yōu)門rue。底層鏈路從任務執(zhí)行到結(jié)果落庫結(jié)果任務的標記最終體現(xiàn)在任務執(zhí)行時的 XCom 推送環(huán)節(jié)。Task SDK 執(zhí)行器在推送返回值 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, )對應到數(shù)據(jù)庫層面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 執(zhí)行 API 一側(cè)POST /xcoms也新增了dag_result布爾查詢參數(shù)見 airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py與任務運行時的推送保持一致。通過 wait API 獲取結(jié)果接口形態(tài)GET /dags/{dag_id}/dagRuns/{dag_run_id}/wait是一個實驗性端點在 OpenAPI 規(guī)范中被標記為experimental說明可能在沒有預告的情況下變更或移除見 v2-rest-api-generated.yaml。其參數(shù)如下參數(shù)位置必填說明dag_idpath是DAG 標識dag_run_idpath是DagRun 標識intervalquery是輪詢 DagRun 狀態(tài)的間隔秒數(shù)必須大于 0exclusiveMinimum: 0.0resultquery否指定要收集結(jié)果 XCom 的任務 id可重復設(shè)置未設(shè)置時默認返回 DAG 中聲明的結(jié)果任務即result或dag返回的 XComArg返回值流式 NDJSON 響應成功響應以換行分隔的 JSONNDJSON流式返回每一行是一個 JSON 對象代表 DagRun 的當前狀態(tài)。OpenAPI 中的響應示例{state: running} {state: success, results: {op: 42}}即運行未結(jié)束時只返回state運行結(jié)束后附帶results字段鍵為任務 id值為該任務返回值 XCom 的內(nèi)容。服務端實現(xiàn)核心邏輯位于DagRunWaiter類airflow-core/src/airflow/api_fastapi/core_api/services/public/dag_run.py其wait方法按interval秒輪詢 DagRun 狀態(tài)并持續(xù)產(chǎn)出 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結(jié)果收集的關(guān)鍵在_serialize_xcoms當result_task_ids is None調(diào)用方未顯式傳result時查詢該 DagRun 全部XCOM_RETURN_KEY且dag_result.is_(True)的 XCom即返回 DAG 作者聲明的結(jié)果任務當調(diào)用方顯式傳了result時按指定的task_ids精確過濾結(jié)果統(tǒng)一按task_id, map_index排序以保證 mapped 任務的結(jié)果順序穩(wěn)定執(zhí)行順序本身不保證非 mapped 任務若只有一條 XCom 則解包為單個值mapped 任務則聚合成按map_index排序的列表。路由處理器wait_dag_run_until_finishedairflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py負責參數(shù)解析與權(quán)限校驗若 DagRun 不存在返回 404若用戶無 XCom 讀取權(quán)限且未顯式請求結(jié)果時會靜默降級為不返回任何 XCom 結(jié)果result_task_ids []若顯式請求了result但無權(quán)限則返回 403。三種調(diào)用方式對照場景請求行為不傳resultDAG 聲明了結(jié)果任務GET /wait?interval1默認返回result/dag返回標記的任務結(jié)果不傳resultDAG 未聲明結(jié)果任務GET /wait?interval1僅返回state無results字段顯式指定任務GET /wait?interval1resulttask_1resulttask_2按指定任務收集可覆蓋 DAG 作者聲明映射任務的聚合行為result與動態(tài)任務映射mapping結(jié)合時所有映射實例的返回值會被聚合成一個按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]}}測試表明單個結(jié)果任務a映射展開后results中a的值為[2, 4]按 map_index 順序即1*2與2*2。這與_serialize_xcoms中_group_xcoms的分組邏輯一致mapped 任務map_index 0的所有 XCom 值以列表返回非 mapped 任務解包為單值。權(quán)限、錯誤與邊界TestWaitDagRun測試類test_dag_run.py系統(tǒng)性地覆蓋了接口的各類邊界401未認證客戶端直接返回 401403無 DAG 訪問權(quán)限返回 403有 RUN 權(quán)限但無 XCOM 權(quán)限時顯式請求result返回 403未顯式請求則降級為不返回結(jié)果狀態(tài)碼仍為 200404DagRun 不存在返回 404422缺少必填的interval參數(shù)返回 422隱式返回值DAG 聲明了結(jié)果任務時{state: ..., results: {task_2: result_2}}DAG 未聲明結(jié)果任務時僅返回{state: ...}顯式返回值?resulttask_1可收集非結(jié)果任務?resulttask_2可收集結(jié)果任務均返回對應任務的返回值 XCom。此外權(quán)限校驗采用了雙重授權(quán)設(shè)計路由依賴先校驗 RUN 訪問權(quán)限處理器內(nèi)再以相同的 team 解析方式校驗 XCOM 訪問權(quán)限避免不同粒度的校驗對 team-aware 認證管理器產(chǎn)生不一致的授權(quán)判斷。小結(jié)result裝飾器與 wait API 構(gòu)成了 Airflow 3.3.0 中結(jié)果感知的 DAG 執(zhí)行閉環(huán)聲明側(cè)result疊加在task之上顯式聲明或通過dag函數(shù)返回XComArg隱式聲明統(tǒng)一落到操作符的returns_dag_result標志執(zhí)行側(cè)Task SDK 推送返回值 XCom 時寫入dag_resultTrue模型與遷移在數(shù)據(jù)庫層面持久化該標記消費側(cè)實驗性 wait 端點以 NDJSON 流式返回 DagRun 狀態(tài)與結(jié)果未顯式指定result參數(shù)時默認返回 DAG 作者聲明的結(jié)果任務mapped 結(jié)果按 map_index 聚合為列表。對于需要以編程方式觸發(fā)并等待 Airflow DAG 完成的調(diào)用方如 CI/CD、數(shù)據(jù)平臺上層編排這提供了一種無需輪詢 XCom 明細即可獲取 DAG最終產(chǎn)出的標準化方式。想深入驗證行為可閱讀 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創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考