authoring-go-sdk-tasks
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`, the `bundlev1` Registry/Dag interfaces, registering Go tasks (`AddTask`/`AddTaskWithName`), dependency injection by parameter type (`context.Context`, `sdk.TIRunContext`, `*slog.Logger`, `sdk.Client`), or reading connections/variables/XComs from Go. This skill covers the Go-specific native API; the shared Python-stub pattern and conceptual model live in authoring-language-sdk-tasks. For building/packing/shipping the bundle see deploying-go-sdk-bundles; for coordinator config see configuring-airflow-language-sdks.
Works with
---
name: authoring-go-sdk-tasks
description: 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`, the `bundlev1` Registry/Dag interfaces, registering Go tasks (`AddTask`/`AddTaskWithName`), dependency injection by parameter type (`context.Context`, `sdk.TIRunContext`, `*slog.Logger`, `sdk.Client`), or reading connections/variables/XComs from Go. This skill covers the Go-specific native API; the shared Python-stub pattern and conceptual model live in authoring-language-sdk-tasks. For building/packing/shipping the bundle see deploying-go-sdk-bundles; for coordinator config see configuring-airflow-language-sdks.
license: Apache-2.0
---
# 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:
```python
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.
```go
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 type | Injected value |
|----------------|----------------|
| `context.Context` | Task context for cancellation. Always available. |
| `sdk.TIRunContext` | Richer context (embeds `context.Context`) exposing `TaskInstance()` and `DagRun()`. See [Runtime context](#runtime-context). |
| `*slog.Logger` | Logger wired to the Airflow task log. |
| `sdk.Client` | Full Airflow model access: Variables, Connections, XComs. |
| `sdk.VariableClient` / `sdk.ConnectionClient` / `sdk.XComClient` | A 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.
```go
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
| Call | Returns | Notes |
|------|---------|-------|
| `GetVariable(ctx, key)` | `(string, error)` | `VariableNotFound` if absent. |
| `UnmarshalJSONVariable(ctx, key, &ptr)` | `error` | Decode 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)` | `error` | Rarely 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.
```go
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.More Data Engineering skills
data-pipeline
claude-office-skills/skills
Data pipeline and ETL automation - extract, transform, load workflows for data integration and analytics
ETL Pipeline
claude-office-skills/skills
Design and automate Extract, Transform, Load data pipelines for data integration and analytics
data-throughput-accelerator
affaan-m/ecc
Use when large data ingestion, backfill, export, ETL, warehouse loading, manifest catch-up, or table synchronization needs to become much faster while preserving data correctness.

