Go SDK
This is an experimental feature.
The Go SDK lets you implement Airflow task logic in Go, with native access to the Airflow “model”
(Variables, Connections, and XCom). The Dag and its scheduling remain in Python; individual tasks delegate
to a compiled Go bundle that is launched by
ExecutableCoordinator for each task instance.
Because Go is a compiled language, every task must be compiled ahead of time and registered inside a single,
self-contained native executable called a bundle. The bundle also embeds the source file of each Dag built with airflow.Dag
plus the entrypoint file (the one with func main), and a metadata
manifest (the dag_id and task_id map) in a footer appended to the executable, so the executable is
the bundle: one runnable file to ship, with no separate manifest or archive. The
airflow-go-pack tool builds and packs that bundle.
API reference
The generated API reference for the Go SDK module, and the list of its released versions, is available on pkg.go.dev.
Prerequisites
Go 1.24 or later to build and pack bundles. This is a build-time requirement only; the worker that runs a packed bundle needs no Go toolchain, because the bundle is a self-contained native executable.
The packed bundle must be accessible from the Airflow worker and the Dag processor, in the Dag bundle the coordinator scans, and built for the operating system and CPU architecture of both.
The
apache-airflow-task-sdkpackage (installed with Airflow) provides the coordinator; no additional Python packages are needed.
Execution architecture
A Python task runner launches the Go bundle directly, with no separate Go worker process on the host. This is the same coordinator mechanism the Java SDK uses. Because the mature Python supervisor handles the Airflow-facing concerns, Go tasks inherit remote task logs (S3/GCS), the full range of task states, and alternate XCom backends rather than implementing them again in Go.
Quick start
The following example shows the minimal moving parts: a Python Dag with two stub tasks, and a Go implementation of those tasks.
Python Dag (the scheduling side)
from airflow.sdk import dag, task
@dag
def simple_dag():
@task.stub(queue="golang")
def extract(): ...
@task.stub(queue="golang")
def transform(): ...
extract() >> transform()
simple_dag()
@task.stub declares the shape of the Go tasks (their names and dependencies) without any Python
implementation. The queue value routes the task to the Go coordinator.
Go implementation
A task is an ordinary Go function whose first parameter is an airflow.Context. Everything Airflow gives
the task is a method on it, so the signature stays the same whatever the task uses.
import (
"runtime"
"github.com/apache/airflow/go-sdk/airflow"
)
func extract(actx airflow.Context) (any, error) {
conn, err := actx.Client().GetConnection(actx, "test_http")
if err != nil {
return nil, err
}
actx.Logger().InfoContext(actx, "fetched connection", "host", conn.Host)
// ... do work, honour actx cancellation ...
return map[string]any{"go_version": runtime.Version()}, nil
}
func transform(actx airflow.Context) error {
val, err := actx.Client().GetVariable(actx, "my_variable")
if err != nil {
return err
}
actx.Logger().InfoContext(actx, "obtained variable", "my_variable", val)
return nil
}
Note
As with the other language SDKs, XCom dependencies are declared in the Python stub Dag (they define task
order). An upstream task’s value reaches a downstream task either through a parameter, when the stub Task
passes it in the TaskFlow call (see Receiving arguments from the stub Task), or by reading it explicitly with
actx.Client().GetXCom.
Go entry point
Build a bundle with airflow.Bundle(), register a handler for each task, and call Serve as the last
statement of main. The Register calls are the single source of truth for which dag_id and task
names this bundle can run, so the generated manifest can never drift from what the binary actually executes.
import (
"log"
"github.com/apache/airflow/go-sdk/airflow"
)
func main() {
bundle := airflow.Bundle()
bundle.Register(
airflow.TaskHandler("simple_dag", "extract", extract),
airflow.TaskHandler("simple_dag", "transform", transform),
)
if err := bundle.Serve(); err != nil {
log.Fatal(err)
}
}
TaskHandler names the dag_id and the task_id explicitly: the dag_id must match the Python
Dag, and the task_id must match a @task.stub function in that Dag. Neither is derived from the Go
function name, so a handler can be named whatever reads best in Go.
TaskHandler also checks the signature of the function it is given and panics if the check fails – for
instance when the function does not take an airflow.Context first, does not return an error, or
declares a variadic ... parameter, which no stub argument can fill. Because main registers every
handler before Serve, a mistake stops the executable as soon as it starts rather than when the task
first runs.
Serve closes registration: a Register left below it in main panics rather than adding a handler
to the map the runtime is already answering from, so what a bundle can run never depends on how far
main has got.
A package that defines task handlers of its own can export them as a []airflow.Registerable for main
to pass on with bundle.Register(reports.Handlers()...).
Coordinator configuration
Register a Dag bundle for the packed bundles, register the coordinator, and route the queue to it in
airflow.cfg (or the equivalent AIRFLOW__* environment variables):
[dag_processor]
dag_bundle_config_list = [
{"name": "dags-folder", "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", "kwargs": {}},
{
"name": "go-task-handlers",
"classpath": "airflow.dag_processing.bundles.local.LocalDagBundle",
"kwargs": {"path": "/opt/airflow/go-task-handlers"}
}
]
[sdk]
coordinators = {
"go": {
"classpath": "airflow.sdk.coordinators.executable.ExecutableCoordinator",
"kwargs": {"task_handler_bundle_name": "go-task-handlers"}
}
}
queue_to_coordinator = {"golang": "go"}
task_handler_bundle_name names the Dag bundle the coordinator scans for packed bundles;
queue_to_coordinator routes stub tasks with queue="golang" to this Go coordinator. See
ExecutableCoordinator configuration for the full list of accepted kwargs and how bundles are located.
There is no separate Go worker to run: the Airflow worker forks the bundle binary once per task instance.
Note
The [sdk] config and the packed bundle files must be present wherever tasks execute and on the Dag
processor. With CeleryExecutor, tasks execute on the Celery workers; with LocalExecutor, they run
inside the scheduler process. The Dag processor checks the stub tasks of each Python Dag against the task
handlers the packed bundles register, so it runs them too and needs bundles built for its operating
system and CPU architecture. The API server does not need any of it. Register the Dag bundle in
[dag_processor] dag_bundle_config_list on every component, like your other Dag bundles: the worker
and the Dag processor resolve task_handler_bundle_name through it, and wherever the [sdk] config
is read it is rejected if the name is missing there.
A Dag processor with this [sdk] configuration also parses the bundle binaries of every Dag bundle (see
Parsing native Dags). Keep it from parsing the binaries that only register task handlers with
that Dag bundle’s .airflowignore, whose pattern syntax follows [core] dag_ignore_file_syntax. Tasks do
not read .airflowignore, so the binaries still run.
With task_handler_bundle_name set to its own Dag bundle, as above, go-task-handlers holds only handler
binaries, so ignore everything in it:
# [core] dag_ignore_file_syntax = glob (the default)
echo '*' > /opt/airflow/go-task-handlers/.airflowignore
# [core] dag_ignore_file_syntax = regexp; a bare * is not a valid pattern and is dropped
echo '.' > /opt/airflow/go-task-handlers/.airflowignore
With task_handler_bundle_name unset, the binaries sit in the same Dag bundle as the Python Dag file. Go
binaries have no extension, so put them in a folder such as bin/ and ignore that folder:
# [core] dag_ignore_file_syntax = glob (the default)
echo 'bin/*' >> /opt/airflow/dags/.airflowignore
# [core] dag_ignore_file_syntax = regexp
echo '^bin/' >> /opt/airflow/dags/.airflowignore
Writing tasks
Every task function takes an airflow.Context as its first parameter, and reaches what Airflow provides
through its methods:
Method |
What it returns |
|---|---|
|
An |
|
A client for Airflow Variables, Connections, and XCom. See The sdk.Client surface. |
|
The identifiers of the running task instance. See Reading the task runtime context. |
|
The identifiers and scheduling timestamps of its Dag run. See Reading the task runtime context. |
airflow.Context is itself a context.Context, so pass it straight to a client call or to
http.NewRequestWithContext, and select on actx.Done(), which fires when the supervisor asks the task
to stop. Respect it for long-running work. Cleanup that must outlive that cancellation runs under context.WithoutCancel(actx).
A helper typed as a plain context.Context recovers the same surface with airflow.FromContext.
Every parameter after the Context is data, filled from the stub Task’s TaskFlow call; see Receiving arguments from the stub Task.
An optional (any, error) return value becomes the task’s return_value XCom. A non-nil error (or a
panic, which the runtime recovers) marks the task instance failed in Airflow, triggering retries if
configured on the stub.
airflow.NewContext builds a Context, so a task is an ordinary function call in a unit test:
actx := airflow.NewContext(
t.Context(), slog.Default(), fakeClient,
airflow.TaskInstance{DagID: "simple_dag", TaskID: "transform", TryNumber: 1},
airflow.DagRun{DagID: "simple_dag", RunID: "run1"},
)
require.NoError(t, transform(actx))
A helper the task calls can still ask for the narrowest interface it needs (for example
sdk.VariableClient instead of the full sdk.Client), which documents the Airflow features it touches
and lets a test pass a fake.
The sdk.Client surface
actx.Client() returns an sdk.Client, which composes four smaller interfaces, so a helper can depend
on just one:
VariableClient-GetVariable(returns the Variable as a string),UnmarshalJSONVariable(decodes a JSON Variable into a pointer you provide),SetVariable, andDeleteVariable.ConnectionClient-GetConnection, returning aConnectionwith fieldsID,Type,Host,Port,Login,Password,Path,Extra(amap[string]any), plus aGetURI()helper.XComClient-GetXComto read an upstream task’s XCom andPushXComto publish one.TaskStateStoreClient-TaskStateStore, returning the store for this task instance:Get,UnmarshalJSONValue(decodes a JSON value into a pointer you provide),Set,Delete, andClear. See The task state store.
GetXCom returns the stored value as an any; see XCom type mapping for how the stored JSON maps to
Go types.
SetVariable stores the value as a string, so encode structured data (for example with json.Marshal)
before storing it.
client := actx.Client()
if err := client.SetVariable(actx, "process_threshold", "42", "Rows above this count take the slow path"); err != nil {
return err
}
if err := client.DeleteVariable(actx, "legacy_threshold"); err != nil {
return err
}
Note
A value supplied by a secrets backend (for example an AIRFLOW_VAR_* environment variable) still takes
precedence over the stored value when the Variable is read back. Calling SetVariable with an empty
description clears any existing description.
Not-found lookups return sentinel errors - VariableNotFound, ConnectionNotFound, XComNotFound,
TaskStateNotFound - so you can branch on a missing value with errors.Is rather than parsing an error
string.
The task state store
actx.Client().TaskStateStore() returns a persistent key/value store private to one task instance,
and the Go SDK’s entry point to durable execution. It is the same store the Python SDK exposes as
context["task_state_store"]; see Task State Store for the concept and its
configuration.
The store is scoped to dag_id, run_id, task_id, and map_index. It deliberately does not
include try_number, so a value written by one attempt is still readable by the next one: a task that
records an external job ID or its own progress can resume after a worker crash or a retry instead of
redoing the work. The Execution API confines every call to the task instance the caller is running as, so
there is no way to address another task’s store - pass results between tasks with XCom instead.
The usual shape is to look for a checkpoint first and only do the expensive work when it is missing:
import (
"errors"
"github.com/apache/airflow/go-sdk/airflow"
"github.com/apache/airflow/go-sdk/sdk"
)
func runSparkJob(actx airflow.Context) error {
store := actx.Client().TaskStateStore()
var jobID string
stored, err := store.Get(actx, "job_id")
switch {
case errors.Is(err, sdk.TaskStateNotFound):
// First attempt: submit the job and remember its ID before doing anything else.
if jobID, err = sparkClient.SubmitJob(actx); err != nil {
return err
}
if err := store.Set(actx, "job_id", jobID, sdk.WithRetention(sdk.NeverExpire)); err != nil {
return err
}
case err != nil:
return err
default:
// Get returns an any; the value was stored by this task as a string.
jobID = stored.(string)
actx.Logger().InfoContext(actx, "reattaching to job submitted by an earlier attempt", "job_id", jobID)
}
return sparkClient.WaitForCompletion(actx, jobID)
}
value must not be nil and must be JSON-representable - a string, number, bool, slice, map, or a struct,
which is stored as an object built from its exported fields and their json tags. A custom
MarshalJSON is not called, so a type that relies on one is stored as the shape of its fields, and a
struct with no exported fields is stored as {}. Read a scalar back with Get, which returns it as an any (the
numeric caveat in XCom type mapping applies here too); for an object or array,
UnmarshalJSONValue decodes it straight into a pointer you provide.
A value the store cannot hold is rejected before it is sent, so you get an error naming the problem rather
than a round trip that fails on the server. The one that catches people out is time.Time, which JSON has
no spelling for - store value.Format(time.RFC3339) and parse it back with time.Parse. Non-finite
floats and []byte are refused for the same reason. This mirrors the Python SDK, where the same values
fail Pydantic validation before the write leaves the worker.
Keys expire, so retention is part of writing a value:
Setwithout options uses the deployment’s[state_store] default_retention_days(30 days by default). The Go runtime cannot read Airflow’s config, so the supervisor resolves that value and passes it in the environment when it launches the bundle. A deployment that sets it to something unusable - a negative number, or a value that is not a whole number of days - fails the write, exactly as it does for a Python task, rather than quietly substituting a different lifetime.sdk.WithRetentiontakes an explicit, positivetime.Duration, orsdk.NeverExpirefor a key that is skipped by garbage collection entirely. A zero or negative retention is rejected rather than given a meaning of its own: to follow the deployment default omit the option, and to drop a key callDelete.
Delete removes one key (deleting a key that does not exist is not an error) and Clear removes
every key stored for this task instance.
Note
The Go SDK does not implement the worker-side state backend ([workers] state_store_backend), which
offloads large values to external storage and records only a reference marker in the database. If a
deployment configures one, a Go task reading a key that was written through that backend receives the raw
reference marker rather than the original value, and a Go task writing a key stores the whole value in the
database instead of offloading it. This is the same behaviour as a Python worker that does not have the
backend configured.
Reading the task runtime context
airflow.Context carries the identifiers and scheduling timestamps of the running task instance and its
Dag run – the Go equivalent of the execution context the Python and Java SDKs expose:
func extract(actx airflow.Context) (any, error) {
ti := actx.TaskInstance()
actx.Logger().InfoContext(actx, "running",
"dag_id", ti.DagID,
"run_id", ti.RunID,
"task_id", ti.TaskID,
"try_number", ti.TryNumber,
"logical_date", actx.DagRun().LogicalDate,
)
return nil, nil
}
actx.TaskInstance() returns DagID, RunID, TaskID, MapIndex (nil for an unmapped task),
and TryNumber; actx.DagRun() returns DagID, RunID, and the *time.Time fields
LogicalDate, DataIntervalStart, and DataIntervalEnd (nil when the run has no such value, e.g. a
manual trigger).
Receiving arguments from the stub Task
A stub Task’s arguments reach a Go handler in one of two ways:
Positional binding – each data parameter takes the argument in the same position.
Struct-based (keyword) binding – a sole struct parameter takes the arguments by field name.
Every parameter after the airflow.Context is a data parameter, filled in declaration order from the
arguments of the Python stub Task’s TaskFlow call. A literal in the Dag file (transform("uk", ...))
decodes straight into the parameter; an upstream task’s output (transform(..., extract())) is pulled
from that task’s XCom in the current Dag run. If the argument count does not match, or an argument’s
declared type cannot fill the Go type, the task fails before its body runs.
// The Python stub Task calls transform("uk", extract()).
func transform(actx airflow.Context, country string, extracted map[string]any) error {
actx.Logger().InfoContext(actx, "transforming", "country", country)
return nil
}
When a task’s sole data parameter is a struct, its fields bind by name instead of by position – keyword arguments rather than positional ones. Being the only data parameter is the opt-in; there is no marker to add.
type CombineInput struct {
Region string `arg:"region_code"` // renamed
Threshold float64
}
// The Python stub Task calls combine(region_code="uk", threshold=0.5).
func Combine(actx airflow.Context, input CombineInput) (any, error) {
return nil, nil
}
An exported field binds the argument matching its own Go name, folding case and underscores, so
Threshold takes threshold; reach for an arg:"<name>" tag when the names genuinely differ, as
Region does above. The Go SDK README has the full binding rules, including
how unmatched fields and arguments are treated and when an untagged struct is decoded whole from a single
argument instead.
Stub parameters the Dag author left at their Python defaults are the exception to both shapes: they reach the wire but need no Go parameter, so adding a defaulted parameter to a stub does not break the Go functions already bound to it.
XCom type mapping
XCom values are stored as JSON in Airflow’s metadata database. The table below shows how those JSON types
surface as Go values when read back via GetXCom.
Python type |
JSON |
Go type (from |
|---|---|---|
|
number (integer) |
numeric (see note) |
|
number (decimal) |
|
|
string |
|
|
boolean |
|
|
null |
|
|
array |
|
|
object |
|
Note
GetXCom returns the value exactly as decoded from the transport; there is no typed XCom
deserialization layer yet. The Python supervisor encodes values as msgpack, so a whole number arrives
as a Go integer type (whose width depends on the value) and only a non-integer arrives as float64. Do
not assume a fixed integer width: type-switch over the numeric types you expect, or round-trip the value
through json.Marshal / json.Unmarshal into a typed Go value.
Building and packaging
A plain go build produces a runnable binary, but a deployable bundle (binary + embedded source files +
manifest) must be produced with airflow-go-pack. The packer compiles the bundle and appends the embedded
metadata footer, so the coordinator can read its dag_ids without executing the binary, producing a
single runnable file. The on-disk format the packer emits (the AFBNDL01 footer and the
airflow-metadata.yaml manifest) is the bundle format shared by all native-executable SDKs, specified in
Executable Bundle Spec.
airflow-go-pack ships via the Go 1.24 tool directive, so there is no global install: add
tool github.com/apache/airflow/go-sdk/cmd/airflow-go-pack
to your bundle module’s go.mod and run it with go tool airflow-go-pack. This pins the packer version
per project.
Build and pack in one step; any flags after -- are forwarded verbatim to go build:
go tool airflow-go-pack ./example/bundle -- -trimpath -tags=prod
Use --output <path> to write the packed bundle straight into the directory of the Dag bundle the
coordinator scans (see Deploying):
go tool airflow-go-pack --output /opt/airflow/go-task-handlers/sample-dag-bundle ./example/bundle
Cross-platform builds
The worker and the Dag processor that run a bundle often use a different operating system or CPU
architecture than your build machine (for example, deploying to a Linux host from an Apple-silicon
darwin/arm64 laptop). Pass --goos / --goarch and the packer cross-builds for you:
go tool airflow-go-pack --goos linux --goarch amd64 \
--output /opt/airflow/go-task-handlers/sample-dag-bundle \
./example/bundle
Alternatively, pack a pre-built binary with --executable / --source. The packer normally execs the
binary with --airflow-metadata to read its manifest, but a cross-compiled binary cannot run on the build
host. In that case, generate the manifest on a machine that can run the binary and feed it to the packer
with --airflow-metadata:
# On a linux/amd64 machine:
go build -o my-bundle ./example/bundle
./my-bundle --airflow-metadata > airflow-metadata.yaml
# Back on the darwin/arm64 machine:
go tool airflow-go-pack --executable ./my-bundle --source main.go \
--airflow-metadata airflow-metadata.yaml
(--executable is mutually exclusive with --goos / --goarch and with go build flags after
--, since it packs an already-built binary instead of building one.)
Deploying
Copy or mount the packed bundle into the Dag bundle named by the coordinator’s task_handler_bundle_name.
The ExecutableCoordinator scans that Dag bundle recursively,
matches the incoming dag_id against each bundle’s manifest, verifies the bundle’s integrity hash, and
launches the matching bundle. This scan identifies bundles by the trailer magic, not by file name. The Dag
processor also requires a file name without an extension, see Parsing native Dags.
The matching bundle is marked executable before it is launched, so any Dag bundle works, including an
object-store one such as S3DagBundle that has no concept of file permissions and so cannot preserve
the execute bit the build produced.
Parsing native Dags
Once an ExecutableCoordinator is configured, the Dag processor
parses every bundle binary in each Dag bundle and runs it to collect the Dags it defines. The Dag processor
recognizes a bundle binary only when its file name has no extension, such as bin/orders, not
bin/orders.bin. A binary that only registers task handlers is parsed too: each parse runs it and finds
no Dags, and a binary built with a Go SDK that cannot answer the Dag-parse request records an import error.
Keep those out of the Dag processor with .airflowignore, as described in Quick start.
With one ExecutableCoordinator, it parses the binaries of every Dag
bundle. With several, for example a second one for another task_handler_bundle_name, map each Dag bundle that
holds native Go Dags to one of them in [sdk] dag_bundle_to_coordinator. A binary in a Dag bundle that has no
entry, or whose entry names no ExecutableCoordinator, fails to parse
with an import error:
[sdk]
dag_bundle_to_coordinator = {"dags-folder": "go"}
The Code view shows each native Dag’s own source file, taken from the sources the bundle embeds.
A task of a native Go Dag runs the bundle binary its Dag was parsed from.
ExecutableCoordinator configuration
All kwargs in the coordinators config entry are passed to the
ExecutableCoordinator constructor:
Parameter |
Default |
Description |
|---|---|---|
|
(task’s own Dag bundle) |
Name of the Dag bundle scanned recursively for executable bundles. It is used only by
mixed-language Dags, to locate the task handlers for the |
|
|
Seconds to wait for the bundle subprocess to connect after launch. Increase this if your bundle startup is slow (e.g. on constrained hardware). |
Note
Locating bundles. The packed bundles for the @task.stub tasks of a Python Dag are read from a Dag
bundle, so they are delivered, refreshed and versioned by the same machinery as your Dags.
The expected layout is a separate Dag bundle for the packed bundles, named by
task_handler_bundle_name, rather than the Dag bundle that holds your.pyfiles. The task uses the version that Dag bundle is on when it starts, pinned for the whole task.If
task_handler_bundle_nameis unset, packed bundles are read from the task’s own Dag bundle, pinned to the version the run was created with.
Limitations
A Python stub Dag is still required. The Execution API does not yet carry Dag structure for non-Python languages, so task names and dependencies are declared in Python with
@task.stub. This is a documented known limitation.