authoring-language-sdk-tasks

作者: astronomer

Airflow語言SDK的語言中立基礎——在DAG保留於Python的同時,以非Python語言實作任務邏輯。當使用者…

npx skills add https://github.com/astronomer/agents --skill authoring-language-sdk-tasks

Authoring Language SDK Tasks (Shared Foundation)

Airflow language SDKs let you implement task logic in a language other than Python while the DAG and its scheduling stay in Python. This skill describes the parts that are identical across every language SDK. Each language has its own companion skill for the native API, build tooling, and runtime — see Per-language skills.

Experimental. The language SDKs are in preview. APIs and artifact coordinates may change.


The model

A DAG is authored in Python as usual. Tasks that should run in another language are declared as stubs routed to a dedicated queue. At runtime, Airflow hands a stub task to a coordinator that launches a short-lived native subprocess for that one task instance, runs your compiled/native code, and shuts the subprocess down.

Consequences that hold for every language SDK:

  • One subprocess per task instance — there is no shared in-process state between task instances. Pass data via XCom or an external store.
  • The DAG, schedule, retries, and queue routing live in Python. The native side only implements task logic.
  • Data crossing the boundary is JSON. See The XCom-as-JSON contract.

The two-sided model

Every task has two halves that must agree:

  1. A Python stub in a normal DAG file — no logic; it declares the task, its queue, the dependency graph, and retry policy.
  2. A native implementation (Java, Go, etc.) whose IDs match the Python side and where the work happens.

Python side (scheduling)

The example below uses the Go SDK to be concrete, but the Python side is identical for every language SDK. The queue name ("golang" here) is an arbitrary label you choose — it just has to match a key in queue_to_coordinator (see configuring-airflow-language-sdks). Pick whatever name fits the SDK you're routing to.

from datetime import timedelta
from airflow.sdk import dag, task


@dag
def sales_pipeline():
    @task.stub(queue="golang")          # queue selects the coordinator (see configuring-airflow-language-sdks)
    def extract(): ...

    @task.stub(queue="golang")
    def transform(extracted): ...        # arg only declares the dependency

    @task.stub(queue="golang", retries=1, retry_delay=timedelta(seconds=5))
    def load(transformed): ...

    @task()                              # an ordinary Python task can sit downstream
    def report(loaded):
        print(f"done: {loaded}")

    report(load(transform(extract())))


sales_pipeline()

Rules that apply regardless of language:

  • The stub function name is the task ID and the @dag name (or dag_id=) is the DAG ID. The native side must use these exact IDs.
  • An upstream argument on a stub (e.g. transform(extracted)) exists only to declare the dependency in Python. The value itself is fetched on the native side via XCom — passing it in Python does not hand it to the native code.
  • Queue, retries, and other task arguments are set on the stub, not in the native code. A native task that fails is reported back to Airflow, which then applies the stub's retry policy.
  • The queue value is what routes the task to a coordinator; the same string must appear in queue_to_coordinator (see configuring-airflow-language-sdks).

The XCom-as-JSON contract

XCom values are stored as JSON in Airflow's metadata database, so the boundary between Python and any native language is JSON. The Python/JSON side is the same for every SDK:

Python typeJSON
intnumber (integer)
floatnumber (decimal)
strstring
boolboolean
Nonenull
listarray
dictobject

Each language SDK maps these JSON types onto its own native types (e.g. a JSON integer becomes a Java Long). The native-type mapping lives in that language's skill. The key portability rule: a value pushed by one task is read by another as JSON, so the consuming side must expect a type compatible with what was stored.


What is language-specific (and lives elsewhere)

This skill deliberately stops at the shared concepts. The following differ per language and are documented in each language's companion skills:

  • Native task API — how you declare tasks, read connections/variables/XComs, and push results (annotations, interfaces, function registration, etc.).
  • Native type mapping — the native column of the JSON table above.
  • Build and packaging — how the artifact is compiled and bundled.
  • Runtime prerequisite — what must be present on the worker (a language runtime for some SDKs, e.g. a JRE for the Java SDK; none for the Go SDK's self-contained bundles).

The Airflow-side wiring (which coordinator runs which queue) is shared in structure but has per-coordinator options; it lives in configuring-airflow-language-sdks.


Language-agnostic pitfalls

  • IDs must match exactly across the Python stub function name and the native task ID, and across @dag/dag_id and the native DAG ID. Mismatches surface as "no DAGs" or missing-XCom errors.
  • Both sides need the upstream reference. Python declares the dependency by passing the upstream call; the native code retrieves the value via XCom.
  • Set queue and retries on the stub, never in the native code.
  • Stub bodies must be empty. An AST check enforces it — only pass, ..., or a docstring is allowed in the body; any real logic is rejected.
  • retry_policy is rejected on stubs (@task.stub raises ValueError). Use retries/retry_delay instead — a retry-policy callable runs Python in-process and would never fire for a task executing in a native subprocess.
  • Assets, deferral, and some other Airflow features have limited or no support in the language SDKs today.

Per-language skills

  • authoring-java-sdk-tasks: Java/Kotlin/JVM native API, type mapping, and logging.
  • authoring-go-sdk-tasks: Go native API — task registration, dependency injection by parameter type, and client access.
  • (Future language SDKs each add their own authoring-<lang>-sdk-tasks skill that builds on this one.)

Related Skills

  • configuring-airflow-language-sdks: Route a queue to a coordinator and set runtime options.
  • authoring-dags: General Airflow DAG authoring (the Python side lives here too).
  • deploying-java-sdk-bundles: Build and ship the Java artifact.
  • deploying-go-sdk-bundles: Build, pack, and ship the Go bundle (per-language deploy skills follow the same shape).

來自 astronomer 的更多技能

airflow-state-store
astronomer
Persists task and asset state across retries and DAG runs using Airflow 3.3's AIP-103 key/value stores (`task_state_store`, `asset_state_store`) and the…
creating-openlineage-extractors
astronomer
為不支援的Airflow運算子及複雜血緣場景設計的自訂OpenLineage提取器。提供兩種方法:直接在你擁有的運算子中加入OpenLineage方法(建議做法),或為無法修改的第三方運算子建立自訂提取器。提取器在三個時間點攔截運算子執行:執行前取得靜態血緣、成功後取得執行階段決定的輸出、以及選擇性地在失敗後取得部分血緣。可透過airflow.cfg或環境變數註冊提取器...
debugging-dags
astronomer
針對失敗的 Airflow DAG 進行系統性根本原因分析與修復,並提供結構化的調查流程。引導完成四個階段的診斷步驟:識別失敗、提取錯誤細節、收集背景資訊,以及提供可行的修復步驟。將失敗分為四種類型(資料、程式碼、基礎設施、相依性),以聚焦調查並建議適當的修正方式。提供可直接使用的 CLI 指令,用於日誌擷取、執行比較、任務清除與 DAG...
delegating-to-otto
astronomer
驅動 Astronomer 的 Otto 代理
deploying-airflow
astronomer
部署 Airflow DAG 和專案。當使用者想要部署程式碼、推送 DAG、設定 CI/CD、部署到生產環境,或詢問部署策略時使用…
deploying-go-sdk-bundles
astronomer
建置、打包並部署已編譯的 Airflow Go SDK 套件,以便 ExecutableCoordinator 能執行它們。當使用者想要編譯 Go 任務套件、要求…時使用。
testing-dags
astronomer
針對Airflow DAG的反覆測試-除錯-修復循環,提供全面的失敗診斷。從af runs trigger-wait <dag_id>開始執行DAG並等待完成,無需預先檢查。失敗時,使用af runs diagnose獲取完整的失敗摘要,並透過af tasks logs檢查特定任務的錯誤細節。支援自訂配置、超時設定與重試機制;能處理成功、失敗及超時情境,並提供清晰的回應解讀。快速驗證功能亦已就緒...
tracing-downstream-lineage
astronomer
追蹤下游資料血緣,在修改資料表或DAG前評估變更影響。透過原始碼搜尋、檢視相依性及BI工具連線,識別目標資料表或DAG的直接消費者。建立完整的相依性樹狀圖,繪製從資料表到儀表板再到機器學習模型的所有下游影響。依關鍵性(關鍵、高、中、低)分類相依性,以優先處理利害關係人溝通與測試。產出包含風險評估、受影響範圍的影響報告。