airflow-new-sdk

作成者: astronomer

Airflow(AIP-108)向けに全く新しい言語SDKを実装するためのガイドです。新しいプログラミング言語のサポートを追加したいコントリビューターは、このスキルを使用してください。…

npx skills add https://github.com/astronomer/airflow --skill airflow-new-sdk

Implementing a new language SDK for Airflow

Start here

Read contributing-docs/30_new_language_sdk.rst first. It is the authoritative contributor guide for this topic — coordinator base class choices, wire protocol spec, bundle footer format, and testing requirements. Everything in this skill builds on top of it, not alongside it.


Repository layout

Every new SDK needs two things. The coordinator (Python) goes here:

task-sdk/src/airflow/sdk/coordinators/<language>/
    __init__.py        # re-export + module docstring
    coordinator.py     # subclass of SubprocessCoordinator or BaseCoordinator
task-sdk/tests/coordinators/<language>/
    test_coordinator.py
task-sdk/tests/integration/coordinators/<language>/
    test_integration.py   # requires Breeze

The language SDK itself lives in a top-level <language>-sdk/ directory (like java-sdk/ and go-sdk/). For native-executable languages using ExecutableCoordinator, no coordinator code is needed at all.


Choosing the right base class — quick guide

Does the runtime compile to a self-contained native executable?
  YES → Use ExecutableCoordinator (zero Python to write).
        Append an AFBNDL01 footer with a packer tool (see go-sdk reference).
  NO  →
    Does it start via a single shell command (node, ruby, dotnet, …)?
      YES → Subclass SubprocessCoordinator.
            Implement _build_execute_task_command only (see 30_new_language_sdk.rst).
      NO  →
        Subclass BaseCoordinator and implement execute_task from scratch.
        (Rare: gRPC daemons, shared memory, persistent processes.)

The full rationale, method signature, and socket lifecycle for each path are in 30_new_language_sdk.rst. Read that section before writing any code.


Reference implementations to study

What to studyWhere
SubprocessCoordinator base classtask-sdk/src/airflow/sdk/coordinators/_subprocess.py
Java coordinator (SubprocessCoordinator subclass)task-sdk/src/airflow/sdk/coordinators/java/coordinator.py
ExecutableCoordinator (native bundles)task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py
Wire protocol in Kotlinjava-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/
Wire protocol in Gogo-sdk/pkg/execution/
AFBNDL01 footer (Go reference)go-sdk/internal/bundlefooter/, task-sdk/docs/executable-bundle-spec.rst
All message types and field specstask-sdk/src/airflow/sdk/execution_time/schema/schema.json

The Java and Go implementations are the two production reference points. When implementing a new SDK, read whichever matches the target language's runtime model (JVM/interpreted → Java; native/compiled → Go).


Logging

The Logging section of 30_new_language_sdk.rst is the spec: the --logs JSON record format, the level names, and the AIRFLOW__LOGGING__* environment variables. Read it first. A few language-neutral details it leaves out:

  • Level values. Levels follow Python's logging scale, so thresholds line up with the rest of Airflow: CRITICAL=50, ERROR=40, WARNING=30, INFO=20, DEBUG=10, NOTSET=0.
  • Parsing NAMESPACE_LEVELS. Split the value on [\s,]+, then split each item on = into (logger_name, level_name). Emit a record only when its level is >= the matching logger_name threshold, or the global threshold when no per-logger entry matches.
  • Don't drop late logs. Connect the --logs socket early and keep it open until the --comm channel has finished; otherwise records emitted during teardown can be lost.
  • Extra config keys. The runtime can't read Airflow's config, so if your SDK needs [logging] settings beyond the two above, propagate them as environment variables from your coordinator's start, the same way.

For how a given language wires its native logging frameworks into this channel, read that SDK's source (e.g. java-sdk/) rather than reproducing it here.

E2E test suite

Create two files mirroring java_sdk_tests/ or go_sdk_tests/:

airflow-e2e-tests/tests/airflow_e2e_tests/<language>_sdk_tests/
    __init__.py
    test_<language>_sdk_dag.py

The test file should:

  • Trigger the SDK's example Dag (the one added to <language>-sdk/dags/ or equivalent) via AirflowClient.trigger_dag.
  • Wait for the run to finish with AirflowClient.wait_for_dag_run.
  • Assert that each SDK task instance reached "success".
  • Assert at least one XCom value — confirms the full round-trip from task return value through the supervisor to the XCom store.
  • Assert structured logs where the SDK emits them (see the Go suite for an example of log-content assertions).

Run locally with:

E2E_TEST_MODE=<language>_sdk uv run --project airflow-e2e-tests pytest \
    tests/airflow_e2e_tests/<language>_sdk_tests/ -xvs

The Java (java_sdk_tests/test_java_sdk_dag.py) and Go (go_sdk_tests/test_go_sdk_dag.py) suites are the reference implementations.


PR checklist (items not covered by 30_new_language_sdk.rst)

The 30_new_language_sdk.rst guide covers coordinator placement, wire protocol implementation, and testing. These additional items belong in the same PR:

  1. task-sdk/src/airflow/sdk/coordinators/<language>/__init__.py — short module docstring and __all__ re-export of the coordinator class.
  2. airflow-core/docs/authoring-and-scheduling/language-sdks/<language>.rst — user-facing doc following the structure of java.rst or go.rst.
  3. airflow-core/docs/authoring-and-scheduling/language-sdks/index.rst — add the new doc to the toctree.
  4. airflow-core/newsfragments/<PR>.feature.rst — a new language is always user-visible; add a newsfragment.
  5. CI wiring — check dev/breeze/src/airflow_breeze/utils/selective_checks.py to confirm <language>-sdk/ changes trigger the right test group. Add if missing.
  6. E2E tests — add a <language>_sdk_tests/ suite under airflow-e2e-tests/tests/airflow_e2e_tests/. See below.

astronomerのその他のスキル

airflow
astronomer
Apache AirflowのDAG、実行、タスク、システム設定をクエリ、管理、トラブルシューティングします。DAG検査、実行管理、タスクログ、設定クエリ、REST API直接アクセスを含む30以上のコマンドをサポート。複数のAirflowインスタンスを永続的な設定で管理し、ローカルおよびAstroデプロイメントを自動検出。DAG実行を同期的(完了待機)または非同期的にトリガーし、障害を診断、再試行のために実行をクリア、リトライ/マップインデックスフィルタリング付きでタスクログにアクセス。出力...
official
airflow-hitl
astronomer
人間による承認ゲート、フォーム入力、およびAirflow DAG内での分岐を、遅延可能オペレーターを使用して実現。4種類のオペレーター:承認/却下の判断を行うApprovalOperator、フォームによる複数選択肢の選択を行うHITLOperator、人間主導のタスクルーティングを行うHITLBranchOperator、フォームデータ収集を行うHITLEntryOperator。すべてのオペレーターは遅延可能であり、Airflow UIのRequired Actionsタブまたは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によるテーブルスキーマ検出を含みます。分析用にPolarsまたはPandas DataFrameを返すrun_sql()およびrun_sql_pandas()カーネル関数を提供します。概念、パターン、テーブルキャッシュを管理するCLIコマンド、さらに...
official
annotating-task-lineage
astronomer
Airflowタスクにデータ系列を注釈付けし、インレットとアウトレットを使用します。OpenLineage Datasetオブジェクト、Airflow Assets、Airflow Datasetsをサポートし、データベース、データウェアハウス、クラウドストレージ間での入出力を定義します。オペレーターに組み込みのOpenLineage抽出機能がない場合のフォールバックとして使用し、カスタム抽出機能とOpenLineageメソッドが優先される4段階の優先順位システムに従います。Snowflake、BigQuery、S3、PostgreSQL向けのデータセット命名ヘルパーを含み、一貫性を確保します。
official
authoring-dags
astronomer
Apache Airflow DAGを作成するためのガイド付きワークフローで、検証とテストの統合を備えています。構造化された6フェーズのアプローチ:環境と既存のパターンを発見し、DAG構造を計画し、ベストプラクティスに従って実装し、af CLIコマンドで検証し、ユーザーの同意を得てテストし、修正を繰り返します。発見用のCLIコマンド(af config connections、af config providers、af dags list)と検証用のCLIコマンド(af dags errors、af dags get、af dags explore)は、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、Kotlin、または任意のJVM言語でAirflow Java SDKを使用して記述します。ユーザーがJava/JVMでAirflowタスクを実装したい場合、または…と尋ねた場合に使用します。
official