)
Apache Airflow 101使用 Airflow SDK 編寫你的第一個 DAG 工作流附完整源碼解析與測試指南【免費下載鏈接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows項目地址: https://gitcode.com/GitHub_Trending/ai/airflow本教程以 Apache Airflow 官方入門文檔 fundamentals.rst 為主線結合倉庫內完整示例 tutorial.py系統講解 DAG 的概念、DAG 定義文件的編寫方式、Operator 與 Task 的關系、Jinja 模板渲染、任務依賴編排以及命令行測試方法。讀完本文你將具備獨立編寫、校驗、本地測試并提交一個可被 Scheduler 調度執行的 Airflow 工作流的能力。什么是 DAGDAGDirected Acyclic Graph有向無環圖是 Airflow 中工作流的核心抽象。簡單來說DAG 是一組任務的集合這些任務按照它們之間的關系與依賴被組織起來——它就像一張工作流的路線圖清楚地展示了每個任務如何與其他任務連接。有向意味著任務之間的依賴有方向誰先誰后無環意味著依賴關系不能形成閉環否則調度器無法確定執行起點。Airflow 會在解析 DAG 時檢測環路一旦發現循環依賴或同一依賴被重復引用就會拋出錯誤見下文設置任務依賴一節。一個完整的 Pipeline 定義示例倉庫中的 tutorial.py 是官方入門示例雖然初看有些內容但每一行都值得拆解。我們先給出完整代碼隨后逐段解釋# [START tutorial] import textwrap from datetime import datetime, timedelta # Operators; we need this to operate! from airflow.providers.standard.operators.bash import BashOperator # The DAG object; well need this to instantiate a DAG from airflow.sdk import DAG with DAG( tutorial, default_args{ depends_on_past: False, retries: 1, retry_delay: timedelta(minutes5), }, descriptionA simple tutorial DAG, scheduletimedelta(days1), start_datedatetime(2021, 1, 1), catchupFalse, tags[example], ) as dag: t1 BashOperator( task_idprint_date, bash_commanddate, ) t2 BashOperator( task_idsleep, depends_on_pastFalse, bash_commandsleep 5, retries3, ) t1.doc_md textwrap.dedent( \ #### Print the current date This task runs date by using the bash_command argument on BashOperator. ... ) dag.doc_md __doc__ templated_command textwrap.dedent( {% for i in range(5) %} echo {{ ds }} echo {{ macros.ds_add(ds, 7)}} {% endfor %} ) t3 BashOperator( task_idtemplated, depends_on_pastFalse, bash_commandtemplated_command, ) t1 [t2, t3] # [END tutorial]理解 DAG 定義文件把 Airflow 的 Python 腳本想象成一個用代碼描述 DAG 結構的配置文件——這一點非常重要你在其中定義的 task 實際運行在另一個環境Scheduler 派發、Worker 執行中因此這個腳本本身不是用來做數據處理的。它的主要職責是定義DAG對象并且必須能夠被快速求值。原因在于Airflow 的 Dag File ProcessorDAG 文件處理器會定期檢查 DAG 文件夾中的每個文件一旦發現變更就重新解析。如果腳本在導入階段執行了重量級操作例如連接數據庫、發起網絡請求會拖慢整個解析過程甚至導致 DAG 無法被識別。因此最佳實踐是DAG 文件中只做定義工作把真正的計算留給任務執行階段。導入模塊與其他 Python 腳本一樣第一步是導入所需庫。tutorial 示例導入了三部分內容import textwrap from datetime import datetime, timedelta # Operators; we need this to operate! from airflow.providers.standard.operators.bash import BashOperator # The DAG object; well need this to instantiate a DAG from airflow.sdk import DAGtextwrap標準庫用于textwrap.dedent去除多行字符串的公共縮進讓文檔字符串與模板字符串書寫更整潔datetime/timedelta標準庫用于設置start_date、retry_delay等時間參數BashOperator來自airflow.providers.standard.operators.bash是本次使用的 OperatorDAG來自airflow.sdk是本倉庫中 DAG 對象的官方入口注意新版 Airflow 已將其收斂到airflow.sdk命名空間。關于 Python 與 Airflow 的模塊管理機制例如哪些目錄會被自動掃描、模塊名如何解析可參考倉庫文檔 modules_management.rst。設置默認參數default_args創建 DAG 及其任務時你可以把參數直接傳給每個 task也可以把一組公共參數放在字典中統一定義。后者通常更高效、更整潔因為所有 task 會默認繼承這些參數同時仍允許在單個 task 上覆蓋。tutorial 示例中的default_argsdefault_args{ depends_on_past: False, retries: 1, retry_delay: timedelta(minutes5), # queue: bash_queue, # pool: backfill, # priority_weight: 10, # end_date: datetime(2016, 1, 1), # wait_for_downstream: False, # execution_timeout: timedelta(seconds300), # on_failure_callback: some_function, # or list of functions # on_success_callback: some_other_function, # or list of functions # on_retry_callback: another_function, # or list of functions # sla_miss_callback: yet_another_function, # or list of functions # on_skipped_callback: another_function, #or list of functions # trigger_rule: all_success },被注釋掉的參數展示了default_args的常見能力邊界它們都可以按需啟用參數含義示例值depends_on_past當前任務實例是否依賴上一個調度周期的任務實例成功Falseretries失敗后自動重試的次數1retry_delay兩次重試之間的等待時長timedelta(minutes5)queue任務被派發到的隊列bash_queuepool任務使用的資源池backfillpriority_weight任務優先級權重10end_date任務不再被調度的時間點datetime(2016, 1, 1)wait_for_downstream是否等待下游任務的上一個實例成功Falseexecution_timeout任務實例運行的最長時限超時則失敗timedelta(seconds300)on_failure_callback/on_success_callback/on_retry_callback/on_skipped_callback失敗/成功/重試/跳過時的回調函數或函數列表some_functionsla_miss_callbackSLA 未達成的回調yet_another_functiontrigger_rule任務觸發條件規則all_success如果想深入了解BaseOperator的全部參數可查閱 airflow.sdk.BaseOperator 文檔 及倉庫中 BaseOperator 源碼 對應的BaseOperator類定義。創建 DAG接下來實例化一個DAG對象來承載任務。tutorial 示例with DAG( tutorial, default_args{...}, descriptionA simple tutorial DAG, scheduletimedelta(days1), start_datedatetime(2021, 1, 1), catchupFalse, tags[example], ) as dag:各參數含義tutorial位置參數dag_id即 DAG 的唯一標識符。它在你的 Airflow 實例中必須全局唯一UI、CLI、API 都通過它引用該 DAGdefault_args上一步定義的默認參數字典會傳遞給每個任務descriptionDAG 的簡要描述展示在 UI 的 DAG 列表中scheduletimedelta(days1)調度周期這里表示每天運行一次。除了timedelta還可以使用 cron 表達式字符串如0 0 * * *或daily等預設值start_dateDAG 開始生效的時間。注意 Airflow 的調度是結束時間對齊的第一個 DAG Run 的 logical date邏輯日期通常等于start_date但真正觸發時間在其之后的一個調度周期catchupFalse關閉回填。若為True當 DAG 從start_date到當前時間之間有大量錯過的調度周期時Scheduler 會一次性補跑所有錯過的 DAG Run設為False則只調度最新周期避免剛上線時批量觸發tags[example]為 DAG 打標簽便于在 UI 中按標簽篩選with ... as dag:上下文管理器語法塊內實例化的任務會自動注冊到該 DAG 上這是新版 Airflow 推薦、也是本倉庫示例采用的寫法。理解 OperatorOperator 是 Airflow 中的工作單元是構建工作流的積木決定了任務將要執行什么動作。所有 Operator 都繼承自BaseOperator因此共享運行任務所需的核心參數如retries、depends_on_past、execution_timeout等。社區與官方提供大量 Operator常見的有PythonOperator執行 Python 可調用對象BashOperator執行 Bash 命令或腳本本教程主角KubernetesPodOperator在 Kubernetes 集群中拉起 Pod 執行任務以及各類 Provider 提供的專有 Operator可瀏覽倉庫 providers 目錄。除了直接實例化 OperatorAirflow 還提供更Pythonic的 TaskFlow API可以用裝飾器把普通 Python 函數變成任務本文暫不展開。定義任務Task要使用 Operator必須先把它實例化為任務Task。任務決定了 Operator 在 DAG 上下文中如何執行工作。task_id是每個任務的唯一標識符。tutorial 示例實例化了兩次BashOperatort1 BashOperator( task_idprint_date, bash_commanddate, ) t2 BashOperator( task_idsleep, depends_on_pastFalse, bash_commandsleep 5, retries3, )注意這里把Operator 專有參數bash_command與BaseOperator繼承來的通用參數retries、depends_on_past混合使用這讓代碼更簡潔。t2還把retries覆蓋為3演示了按任務覆蓋默認值的能力。任務參數的優先級如下顯式傳入的參數如t2的retries3default_args字典中的值如retries1Operator 自身的默認值如果存在。注意每個任務必須包含或繼承task_id和owner兩個參數否則 Airflow 會報錯。幸運的是全新安裝的 Airflow 默認將owner設為airflow因此你通常只需確保設置task_id即可。使用 Jinja 模板渲染Airflow 內置了Jinja 模板引擎讓你能訪問內置變量如{{ ds }}與宏macros來動態生成命令內容。{{ ds }}是最常用的模板變量代表邏輯日期logical date的日期戳格式為YYYY-MM-DD。tutorial 示例中的模板任務templated_command textwrap.dedent( {% for i in range(5) %} echo {{ ds }} echo {{ macros.ds_add(ds, 7)}} {% endfor %} ) t3 BashOperator( task_idtemplated, depends_on_pastFalse, bash_commandtemplated_command, )這段模板包含{% for i in range(5) %}...{% endfor %}Jinja 控制流語句塊循環 5 次{{ ds }}邏輯日期的日期戳變量例如2021-01-01{{ macros.ds_add(ds, 7) }}宏調用ds_add會在給定日期上加上指定天數這里得到2021-01-08。實際渲染后templated任務執行的命令大致是echo 2021-01-01 echo 2021-01-08 echo 2021-01-01 echo 2021-01-08 ...共 5 組關于模板還有幾點實用技巧傳入腳本文件bash_command可以直接傳文件名例如bash_commandtemplated_command.sh把命令邏輯拆到獨立文件中便于組織與維護自定義宏與過濾器可以在 DAG 上定義user_defined_macros和user_defined_filters創建自己的模板變量與過濾器完整變量/宏清單所有可在模板中引用的變量與宏見倉庫 templates-ref.rst。為 DAG 與任務添加文檔Airflow 允許為 DAG 或單個任務附加文檔直接在 UI 中渲染查看DAG 文檔以Markdown渲染在 DAG 詳情頁任務文檔支持純文本、Markdown、reStructuredText、JSON、YAML 等多種格式。當任務文檔使用doc_md時Airflow 渲染常見的 Markdown 特性包括行內代碼、圍欄代碼塊、以math圍欄包裹的公式用KaTeX渲染以及mermaid圍欄中的Mermaid 流程圖。在 tutorial DAG 中print_date任務t1通過doc_md展示了這些能力t1.doc_md textwrap.dedent( \ #### Print the current date This task runs date by using the bash_command argument on BashOperator. In the Task Instance Details page, Airflow renders this documentation from the tasks doc_md field. After this task succeeds, Airflow can run both downstream tasks: sleep and templated. bash dateMath fences are rendered with KaTeX. This tutorial starts one task and then branches into two downstream tasks:1\ \text{upstream task} 2\ \text{downstream tasks} 3\ \text{tasks}The same dependency is shown as a Mermaid diagram: )與此同時DAG 級文檔可以用模塊 docstring 或直接賦值 python dag.doc_md __doc__ # 使用文件開頭的 docstring # 或者直接寫字符串 dag.doc_md This is a documentation placed anywhere 下圖展示了任務doc_md在 UI 的 Task Instance Details 頁面中的渲染效果Markdown、KaTeX 公式與 Mermaid 圖并存實踐建議把文檔緊挨著它所描述的任務編寫如t1.doc_md ...保持文檔與代碼同步演進。設置任務依賴Airflow 中任務之間可以互相依賴。假設有任務t1、t2、t3可以用多種方式表達依賴t1.set_downstream(t2) # 這表示 t2 需要等 t1 成功運行后才能運行 # 等價于 t2.set_upstream(t1) # 也可以使用位移運算符bit shift鏈式表達 t1 t2 # 反向的上游依賴 t2 t1 # 鏈式多依賴位移運算符更簡潔 t1 t2 t3 # 列表形式設置依賴以下寫法效果相同 t1.set_downstream([t2, t3]) t1 [t2, t3] [t2, t3] t1tutorial 示例使用的正是列表形式的位移運算符寫法t1 [t2, t3]它表達t1成功之后t2與t3并行運行的扇形結構。需要警惕的是Airflow 會檢測 DAG 中的環cycle也會檢測同一依賴被重復引用的情況一旦發現就會拋出錯誤。因此在編寫復雜 DAG 時應避免出現t1 t2 t1之類的循環依賴。處理時區創建一個時區感知time zone aware的 DAG很簡單使用 pendulum 庫提供的時區感知日期時間即可例如pendulum.datetime(2021, 1, 1, tzAsia/Shanghai)。務必避免使用標準庫datetime.timezone對象因為它們在 Airflow 場景下存在已知限制無法攜帶 IANA 時區名稱、轉換行為有缺陷等。Airflow 內部的時間處理統一基于 pendulum倉庫 shared/timezones 目錄下的源碼封裝了相關邏輯供深入研究者參考。回顧完整代碼完成上述步驟后你的代碼應當與倉庫中的 tutorial.py 一致。整體結構如下導入模塊textwrap、datetime、BashOperator、DAG定義default_args用with DAG(...)創建 DAG 對象實例化BashOperator得到任務t1、t2、t3為任務與 DAG 添加doc_md文檔定義 Jinja 模板命令用t1 [t2, t3]聲明依賴。測試你的 Pipeline寫完之后就該測試了。第一步確認腳本能通過解析。把代碼保存為tutorial.py放到airflow.cfg中dags_folder指定的 DAG 目錄默認如~/airflow/dags然后運行python ~/airflow/dags/tutorial.py如果腳本無錯誤地運行結束說明你的 DAG 結構定義正確。注意直接執行時with DAG(...)塊內的任務實例化與依賴聲明都會正常完成但不會真的執行任務。命令行元數據校驗進一步用 CLI 命令驗證元數據# 初始化數據庫表 airflow db migrate # 打印所有已激活的 DAG 列表 airflow dags list # 打印 tutorial DAG 中的任務列表 airflow tasks list tutorial # 打印 tutorial DAG 的 graphviz 可視化表示 airflow dags show tutorial其中airflow db migrate會創建/更新元數據庫 schema首次使用 Airflow 時必須執行airflow dags show tutorial需要安裝 graphviz 支持輸出 DAG 結構的可視化描述。測試任務實例與 DAG Run你可以針對指定的**邏輯日期logical date**測試某個任務實例這模擬了 Scheduler 在某個日期時間點上運行你的任務。關于邏輯日期請注意Scheduler 是為某個具體日期時間運行你的任務而不一定是在那個日期時間運行。邏輯日期logical date是 DAG Run 被命名的那個時間戳它通常對應工作流所處理時間周期的結束時刻——或者是手動觸發 DAG Run 的時刻。Airflow 用邏輯日期來組織和跟蹤每次運行你在 UI、日志和代碼中都是通過它引用某次具體執行的。當通過 UI 或 API 觸發 DAG 時你也可以自行提供邏輯日期從而按某個時間點運行工作流。命令格式為# 命令布局: command subcommand [dag_id] [task_id] [(可選) 日期] # 測試 print_date 任務 airflow tasks test tutorial print_date 2015-06-01 # 測試 sleep 任務 airflow tasks test tutorial sleep 2015-06-01還可以查看模板是如何渲染的# 測試 templated 任務 airflow tasks test tutorial templated 2015-06-01這條命令會輸出詳細日志并實際執行你的 bash 命令——你會看到模板循環展開后的真實命令內容。需要記住airflow tasks test在本地運行任務實例日志輸出到 stdout不在數據庫中記錄狀態。它是調試單個任務實例的便捷工具airflow dags test在本地運行整個 DAG Run適合測試完整 DAG。與tasks test不同它會創建真實的 DAG Run 并在元數據庫中記錄任務狀態因此需要已初始化的數據庫且 DAG 能被 Airflow 從你的 DAG 文件夾序列化。更多細節可參考 DAG 調試與 dag.test() 相關章節。接下來做什么到這里你已經成功編寫并測試了第一個 Airflow 工作流。下一步把代碼合并到運行著Scheduler的代碼倉庫中Scheduler 會接管你的 DAG按schedule每天自動觸發執行。進階方向繼續學習 TaskFlow API 教程用更 Pythonic 的方式定義工作流瀏覽 核心概念深入理解 DAG、Task、Operator、Scheduler 等底層機制閱讀 templates-ref.rst掌握全部模板變量與宏參考 模塊管理文檔了解 Airflow 如何加載 Python 模塊。掌握了本教程你就擁有了 Airflow 世界中最重要的一塊基石從會寫腳本到能寫出被調度器可靠執行的工作流。【免費下載鏈接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows項目地址: https://gitcode.com/GitHub_Trending/ai/airflow創作聲明:本文部分內容由AI輔助生成(AIGC),僅供參考