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
astronomer
查询、管理和排查Apache Airflow的DAG、运行记录、任务及系统配置。支持30多种命令,涵盖DAG检查、运行管理、任务日志、配置查询及直接REST API访问。通过持久化配置管理多个Airflow实例;自动发现本地和Astro部署。同步(等待完成)或异步触发DAG运行,诊断故障,清除运行记录以重试,并通过重试/映射索引过滤访问任务日志。输出...
official
airflow-hitl
astronomer
在Airflow DAG中使用可延迟操作符实现人工审批关卡、表单输入和分支。四种操作符类型:用于批准/拒绝决策的ApprovalOperator、带表单的多选项选择HITLOperator、人工驱动的任务路由HITLBranchOperator,以及表单数据收集HITLEntryOperator。所有操作符均为可延迟设计,在通过Airflow UI的"必需操作"标签页或REST API等待人工响应时释放工作槽位。支持包括自定义在内的可选功能...
official
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…
official
analyzing-data
astronomer
查询数据仓库,利用缓存的模式和概念映射来回答业务问题。支持对重复问题类型进行模式查找和缓存,并通过记录结果来改进后续查询。包含概念到表的映射缓存,以及通过INFORMATION_SCHEMA或代码库grep进行表结构发现。提供run_sql()和run_sql_pandas()内核函数,返回Polars或Pandas DataFrame用于分析。提供CLI命令用于管理概念、模式和表缓存,以及...
official
annotating-task-lineage
astronomer
使用入口和出口为Airflow任务标注数据血缘。支持使用OpenLineage Dataset对象、Airflow Assets和Airflow Datasets定义跨数据库、数据仓库及云存储的输入输出。当运算符缺少内置OpenLineage提取器时作为备用方案;遵循四级优先级系统,其中自定义提取器和OpenLineage方法优先。包含针对Snowflake、BigQuery、S3和PostgreSQL的数据集命名辅助工具,以确保一致性...
official
authoring-dags
astronomer
创建Apache Airflow DAG的引导式工作流,集成验证与测试。采用六阶段结构化方法:发现环境与现有模式、规划DAG结构、遵循最佳实践实现、通过af CLI命令验证、经用户同意测试、迭代修复。用于发现(af config connections、af config providers、af dags list)和验证(af dags errors、af dags get、af dags explore)的CLI命令可提供DAG的即时反馈...
official
authoring-go-sdk-tasks
astronomer
Writes Airflow task logic in Go using the Airflow Go SDK. Use when the user wants to implement Airflow tasks in Go, asks about `BundleProvider`/`RegisterDags`,…
official
authoring-java-sdk-tasks
astronomer
使用 Airflow Java SDK 编写 Java、Kotlin 或任何 JVM 语言的 Airflow 任务逻辑。当用户想要在 Java/JVM 中实现 Airflow 任务时使用,询问……
official