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`,…

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

Authoring Go SDK Tasks

The Airflow Go SDK implements the language-SDK model for Go: your DAG stays in Python, and each task is a compiled Go function registered inside a bundle (a single native executable). This skill covers the Go-specific native API. The shared model (the Python @task.stub pattern, ID matching, the XCom-as-JSON contract) lives in authoring-language-sdk-tasks; read that first if you are new to language SDKs.

Experimental. The Go SDK is under active development and not production-ready. Module path github.com/apache/airflow/go-sdk (Go 1.24+). APIs may change.

Related skills: authoring-language-sdk-tasks (shared Python stub + concepts), deploying-go-sdk-bundles (build, pack, and ship the bundle), configuring-airflow-language-sdks (route the queue to the Go coordinator).


Recap: the Python side

A Go task is paired with a Python stub that carries no logic; it declares the task, its queue, and the dependency graph. IDs must match the Go registration exactly, and queue= routes the task to the Go runtime. Full rules are in authoring-language-sdk-tasks; the minimal shape:

from airflow.sdk import dag, task


@task.stub(queue="golang")
def extract(): ...


@task.stub(queue="golang")
def transform(): ...


@dag()
def simple_dag():
    extract() >> transform()


simple_dag()

The queue value ("golang" here) is an arbitrary label that must match the queue routed to the Go coordinator (queue_to_coordinator). See configuring-airflow-language-sdks.


The bundle entry point

A bundle implements bundlev1.BundleProvider: report its version and register your DAGs and tasks. main is one line; bundlev1server.Serve wires the bundle to the Airflow runtime for you.

package main

import (
	"log"

	v1 "github.com/apache/airflow/go-sdk/bundle/bundlev1"
	"github.com/apache/airflow/go-sdk/bundle/bundlev1/bundlev1server"
)

type myBundle struct{}

var _ v1.BundleProvider = (*myBundle)(nil)

func (m *myBundle) GetBundleVersion() v1.BundleInfo {
	return v1.BundleInfo{Name: bundleName, Version: &bundleVersion}
}

func (m *myBundle) RegisterDags(dagbag v1.Registry) error {
	simpleDag := dagbag.AddDag("simple_dag")      // dag_id must match the Python @dag name
	simpleDag.AddTask(extract)                    // task_id is the function name; must match the stub
	simpleDag.AddTaskWithName("transform", transform) // or set the task_id explicitly
	return nil
}

func main() {
	if err := bundlev1server.Serve(&myBundle{}); err != nil {
		log.Fatal(err)
	}
}

AddTask(fn) derives the task_id from the Go function's name; use AddTaskWithName("<task_id>", fn) when that name can't match the Python stub (an unexported, renamed, or reused function). RegisterDags is the single source of truth for task identity: the bundle's manifest (used by the packer and by the coordinator) is generated by running it, never hand-written.


Task functions: dependency injection by parameter type

A task is an ordinary Go function. The runtime inspects its signature and injects arguments by type; declare only what you need.

Parameter typeInjected value
context.ContextTask context for cancellation. Always available.
sdk.TIRunContextRicher context (embeds context.Context) exposing TaskInstance() and DagRun(). See Runtime context.
*slog.LoggerLogger wired to the Airflow task log.
sdk.ClientFull Airflow model access: Variables, Connections, XComs.
sdk.VariableClient / sdk.ConnectionClient / sdk.XComClientA narrower slice of sdk.Client. Prefer the narrowest you need; it documents intent and is trivial to fake in tests.

The optional return signature is (result, error): a non-nil result is pushed as the task's return_value XCom; a non-nil error fails the task (which triggers the stub's retry policy). Returning only error, or nothing, is also valid.

func extract(ctx sdk.TIRunContext, client sdk.Client, log *slog.Logger) (any, error) {
	conn, err := client.GetConnection(ctx, "test_http")
	if err != nil {
		return nil, err
	}
	log.Info("connected", "host", conn.Host)
	return map[string]any{"go_version": runtime.Version()}, nil
}

func transform(ctx sdk.TIRunContext, client sdk.VariableClient) error {
	val, err := client.GetVariable(ctx, "my_variable")
	if err != nil {
		return err // VariableNotFound (a sentinel error) if absent
	}
	_ = val
	return nil
}

The sdk.Client surface

CallReturnsNotes
GetVariable(ctx, key)(string, error)VariableNotFound if absent.
UnmarshalJSONVariable(ctx, key, &ptr)errorDecode a JSON variable into a struct/pointer.
GetConnection(ctx, connID)(Connection, error)ConnectionNotFound if absent.
GetXCom(ctx, dagID, runID, taskID, mapIndex, key, value)(any, error)XComNotFound only if the key is absent; a stored null returns (nil, nil).
PushXCom(ctx, ti, key, value)errorRarely needed; a returned value is pushed for you.

Connection exposes ID, Type, Host, Port (int), Login *string, Password *string (nil when unset, distinct from empty), Path (schema), Extra map[string]any, plus GetURI(). Not-found cases return the sentinels sdk.VariableNotFound, sdk.ConnectionNotFound, sdk.XComNotFound.

To read an upstream task's result, call GetXCom explicitly, taking the dag_id/run_id/task_id you need from the runtime context (below).


Runtime context

Declare an sdk.TIRunContext parameter to read metadata about the task instance and its DAG run. It is an interface that embeds context.Context, so it is usable anywhere a context.Context is expected.

func extract(ctx sdk.TIRunContext, log *slog.Logger) error {
	ti, dagRun := ctx.TaskInstance(), ctx.DagRun()
	log.Info("running",
		"task_id", ti.TaskID,
		"run_id", dagRun.RunID,
		"logical_date", dagRun.LogicalDate)
	return nil
}
  • TaskInstance(): DagID, RunID, TaskID, MapIndex *int (nil when unmapped), TryNumber.
  • DagRun(): DagID, RunID, and the *time.Time timestamps LogicalDate, DataIntervalStart, DataIntervalEnd (nil when not sent).

The accessors are populated from the task's startup details before the body runs. Because TIRunContext embeds context.Context, pass it straight to client calls and cancellation checks (ctx.Done()); declare it as your context parameter by default. In tests, build the argument with sdk.NewTIRunContext(ctx, ti, dagRun) (it panics on a nil ctx).


Go-specific pitfalls

  • IDs must match the Python stub (dag_id from AddDag, task_id from the registered function name), and the stub's queue= must route to the Go coordinator, or the task is never delivered.
  • RegisterDags is authoritative. Do not hand-write the manifest; the packer generates it by running RegisterDags.
  • Ask for the narrowest client interface you need (sdk.VariableClient over sdk.Client) for clearer intent and easier fakes.
  • A non-nil error return fails the task and applies the stub's retries; a recovered panic is also a failure.
  • See authoring-language-sdk-tasks for the language-agnostic pitfalls (one process per task instance, set queue and retries on the stub).

Related Skills

  • authoring-language-sdk-tasks: Shared Python-stub pattern and concepts (read first).
  • deploying-go-sdk-bundles: Build and pack the bundle with go tool airflow-go-pack, then deploy it for the coordinator.
  • configuring-airflow-language-sdks: Route the queue to the Go coordinator (ExecutableCoordinator).
  • authoring-dags: General Airflow DAG authoring.

来自 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-java-sdk-tasks
astronomer
使用 Airflow Java SDK 编写 Java、Kotlin 或任何 JVM 语言的 Airflow 任务逻辑。当用户想要在 Java/JVM 中实现 Airflow 任务时使用,询问……
official
authoring-language-sdk-tasks
astronomer
Airflow语言SDK的语言中立基础——在DAG保留于Python的同时,以非Python语言实现任务逻辑。当用户…
official