diff --git a/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst b/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst index d4ff487404999..7e10a7b6dcaec 100644 --- a/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst +++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst @@ -131,8 +131,9 @@ task declares only the parameters it needs. .. note:: As with the other language SDKs, XCom *dependencies* are declared in the Python stub Dag (they define task - order). The value must still be read explicitly in Go via ``client.GetXCom``, and produced either by the - task's ``(any, error)`` return value or by ``client.PushXCom``. + order). A task reads an upstream value either explicitly via ``client.GetXCom`` or by declaring an + :ref:`XCom-input struct parameter `, and produces one via its return value or + ``client.PushXCom``. Go entry point ~~~~~~~~~~~~~~~ @@ -223,6 +224,9 @@ parameters your task actually needs: - A logger whose output is routed back to the Airflow task log. * - ``sdk.Client`` (or a narrower interface) - A client for Airflow Variables, Connections, and XCom. + * - An ``xcom``-tagged struct + - Upstream task XComs, pulled and decoded into the struct's fields before the task runs. See + :ref:`go-sdk/xcom-inputs`. 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 @@ -252,6 +256,40 @@ Go types. Not-found lookups return sentinel errors - ``VariableNotFound``, ``ConnectionNotFound``, ``XComNotFound`` - so you can branch on a missing value with ``errors.Is`` rather than parsing an error string. +.. _go-sdk/xcom-inputs: + +Receiving upstream XComs as parameters +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Instead of pulling XComs by hand with ``GetXCom``, a task can declare an **XCom-input struct** parameter. +Each exported field carries an ``xcom:""`` tag (or ``xcom:",key="`` to pull a key +other than the default ``return_value``); the runtime pulls that upstream task's XCom and decodes it into the +field before the task runs: + +.. code-block:: go + + type TransformInput struct { + Extracted ExtractResult `xcom:"extract"` // extract's return_value + FromPython string `xcom:"python_task_1"` // an upstream Python task's output + } + + func transform(ctx context.Context, log *slog.Logger, in TransformInput) (TransformResult, error) { + // in.Extracted and in.FromPython are already populated. + return TransformResult{Variable: in.FromPython, Extracted: in.Extracted}, nil + } + +Decoding into a struct is strict: an unknown or renamed key fails the task rather than silently leaving +fields zero. To decode loosely instead (for example, when the producer is an evolving or cross-language +task), type the field as ``map[string]any`` (or ``any``), which accepts any shape and keeps unknown keys: + +.. code-block:: go + + type TransformInput struct { + Extracted map[string]any `xcom:"extract"` // any shape; read it with in.Extracted["go_version"] + } + +The upstream dependency must still be declared in the Python stub Dag so the tasks run in order. + .. _go-sdk/types: XCom type mapping diff --git a/go-sdk/README.md b/go-sdk/README.md index 5ceb1afb4cc4b..c2a6dda865a6d 100644 --- a/go-sdk/README.md +++ b/go-sdk/README.md @@ -102,8 +102,9 @@ func main() { A task is an ordinary Go function. The runtime inspects its signature and injects arguments by type: `context.Context`, `*slog.Logger`, and an `sdk.Client` (or a narrower interface such as -`sdk.VariableClient`). An optional `(any, error)` return becomes the task's XCom; an `error` return marks -the task failed. +`sdk.VariableClient`). A parameter can also be an *XCom-input struct* whose fields are pulled from upstream +tasks (see [Receiving upstream XComs as typed parameters](#receiving-upstream-xcoms-as-typed-parameters)). A +return value becomes the task's XCom; an `error` return marks the task failed. ```go func extract(ctx context.Context, client sdk.Client, log *slog.Logger) (any, error) { @@ -126,6 +127,39 @@ Asking for the narrowest interface a task needs (e.g. `sdk.VariableClient` inste unit testing easier and documents which Airflow features the task touches. `RegisterDags` is the single source of truth for which `dag_id`s and `task_id`s a bundle can run. +### Receiving upstream XComs as typed parameters + +Instead of pulling XComs by hand with `client.GetXCom`, a task can declare an **XCom-input struct** +parameter. Each exported field carries an `xcom:""` tag (or `xcom:",key="` to pull a +key other than the default `return_value`); the runtime pulls that upstream task's XCom and decodes it into +the field before your task runs: + +```go +type TransformInput struct { + Extracted ExtractResult `xcom:"extract"` // extract's return_value + FromPython string `xcom:"python_task_1"` // an upstream Python task's output +} + +func transform(ctx context.Context, log *slog.Logger, in TransformInput) (TransformResult, error) { + // in.Extracted and in.FromPython are already populated. + return TransformResult{Variable: in.FromPython, Extracted: in.Extracted}, nil +} +``` + +Decoding into a struct is strict: an unknown or renamed key fails the task rather than silently leaving +fields zero. To decode loosely instead (for example, when the producer is an evolving or cross-language +task), type the field as `map[string]any` (or `any`), which accepts any shape and keeps unknown keys: + +```go +type TransformInput struct { + Extracted map[string]any `xcom:"extract"` // any shape; read it with in.Extracted["go_version"] +} +``` + +A typed return value (here `TransformResult`) is published as the task's `return_value` XCom for a +downstream task to pull. The upstream dependency must still be declared in the Python stub Dag so the tasks +run in order. + ## Deployment modes A bundle can run in two ways. The same bundle binary works in both; you pick one per deployment: