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