jason810496 commented on code in PR #73528:
URL: https://github.com/apache/airflow/pull/73528#discussion_r4070592736
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -268,30 +295,74 @@ so you can branch on a missing value with ``errors.Is``
rather than parsing an e
Reading the task runtime context
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
-Declare an ``sdk.TIRunContext`` parameter on a task to read 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. It is an interface that embeds ``context.Context``, so the same
``ctx`` drives cancellation and
-client calls. The runtime binds it by type, just like the other injected
parameters:
+``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:
.. code-block:: go
- func extract(ctx sdk.TIRunContext, log *slog.Logger) (any, error) {
- ti := ctx.TaskInstance()
- log.Info("running",
+ 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", ctx.DagRun().LogicalDate,
+ "logical_date", actx.DagRun().LogicalDate,
)
return nil, nil
}
-``ctx.TaskInstance()`` returns ``DagID``, ``RunID``, ``TaskID``, ``MapIndex``
(nil for an unmapped task),
-and ``TryNumber``; ``ctx.DagRun()`` returns ``DagID``, ``RunID``, and the
``*time.Time`` fields
+``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).
+.. _go-sdk/arguments:
+
+Receiving arguments from the stub Dag
+~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
Review Comment:
```suggestion
Receiving arguments from the stub Task
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -200,37 +204,59 @@ There is no separate Go worker to run: the Airflow worker
forks the bundle binar
Writing tasks
-------------
-The runtime inspects a task function's signature and injects arguments by
type, so you only declare the
-parameters your task actually needs:
+Every task function takes an ``airflow.Context`` as its first parameter, and
reaches what Airflow provides
+through its methods:
.. list-table::
:header-rows: 1
:widths: 35 65
- * - Parameter type
- - Injected value
- * - ``sdk.TIRunContext``
- - The task's execution context: the cancellation/deadline signal plus the
task instance identifiers and
- Dag run timestamps. Respect it for long-running work. See
:ref:`go-sdk/runtime-context`.
- * - ``*slog.Logger``
- - 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.
+ * - Method
+ - What it returns
+ * - ``actx.Logger()``
+ - An ``*slog.Logger`` whose output is routed back to the Airflow task log.
+ * - ``actx.Client()``
+ - A client for Airflow Variables, Connections, and XCom. See
:ref:`go-sdk/client`.
+ * - ``actx.TaskInstance()``
+ - The identifiers of the running task instance. See
:ref:`go-sdk/runtime-context`.
+ * - ``actx.DagRun()``
+ - The identifiers and scheduling timestamps of its Dag run. See
:ref:`go-sdk/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``.
Review Comment:
minor nit regarding the formatting.
```suggestion
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``.
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -268,30 +295,74 @@ so you can branch on a missing value with ``errors.Is``
rather than parsing an e
Reading the task runtime context
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
-Declare an ``sdk.TIRunContext`` parameter on a task to read 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. It is an interface that embeds ``context.Context``, so the same
``ctx`` drives cancellation and
-client calls. The runtime binds it by type, just like the other injected
parameters:
+``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:
.. code-block:: go
- func extract(ctx sdk.TIRunContext, log *slog.Logger) (any, error) {
- ti := ctx.TaskInstance()
- log.Info("running",
+ 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", ctx.DagRun().LogicalDate,
+ "logical_date", actx.DagRun().LogicalDate,
)
return nil, nil
}
-``ctx.TaskInstance()`` returns ``DagID``, ``RunID``, ``TaskID``, ``MapIndex``
(nil for an unmapped task),
-and ``TryNumber``; ``ctx.DagRun()`` returns ``DagID``, ``RunID``, and the
``*time.Time`` fields
+``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).
+.. _go-sdk/arguments:
+
+Receiving arguments from the stub Dag
+~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
+
+Every parameter after the ``airflow.Context`` is a **data parameter**, filled
in declaration order from the
+arguments of the Python stub Dag'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.
+
+.. code-block:: go
+
+ // The stub Dag calls transform("uk", extract()).
Review Comment:
```suggestion
// The Python stub Task calls transform("uk", extract()).
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -268,30 +295,74 @@ so you can branch on a missing value with ``errors.Is``
rather than parsing an e
Reading the task runtime context
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
-Declare an ``sdk.TIRunContext`` parameter on a task to read 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. It is an interface that embeds ``context.Context``, so the same
``ctx`` drives cancellation and
-client calls. The runtime binds it by type, just like the other injected
parameters:
+``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:
.. code-block:: go
- func extract(ctx sdk.TIRunContext, log *slog.Logger) (any, error) {
- ti := ctx.TaskInstance()
- log.Info("running",
+ 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", ctx.DagRun().LogicalDate,
+ "logical_date", actx.DagRun().LogicalDate,
)
return nil, nil
}
-``ctx.TaskInstance()`` returns ``DagID``, ``RunID``, ``TaskID``, ``MapIndex``
(nil for an unmapped task),
-and ``TryNumber``; ``ctx.DagRun()`` returns ``DagID``, ``RunID``, and the
``*time.Time`` fields
+``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).
+.. _go-sdk/arguments:
+
+Receiving arguments from the stub Dag
+~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
+
+Every parameter after the ``airflow.Context`` is a **data parameter**, filled
in declaration order from the
+arguments of the Python stub Dag'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.
+
+.. code-block:: go
+
+ // The stub Dag 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
+ }
+
+Stub parameters the Dag author left at their Python defaults are the
exception: 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.
+
+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.
+
+.. code-block:: go
+
+ type CombineInput struct {
+ Region string `arg:"region_code"` // renamed
+ Threshold float64
+ }
+
+ // The stub Dag 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
+<https://github.com/apache/airflow/blob/main/go-sdk/README.md>`__ 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.
+
Review Comment:
```suggestion
Stub parameters the Dag author left at their Python defaults are the
exception: 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.
```
##########
go-sdk/README.md:
##########
@@ -169,37 +167,55 @@ data parameters is rejected at registration. A sole
struct parameter also falls
decoding when it gets exactly one passed argument no field claims, so a task
can still take an
upstream object as a single argument.
-### Reading the task runtime context
+### The task context
+
+`airflow.Context` is the first parameter of every task handler -- the Go
equivalent of the execution
+context the Python and Java SDKs expose. Everything Airflow gives the task is
a method on it:
Review Comment:
Would it be better to not mention the other SDKs in the Go-SDK README.
```suggestion
`airflow.Context` is the first parameter of every task handler. Everything
Airflow gives the task is a method on it:
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -268,30 +295,74 @@ so you can branch on a missing value with ``errors.Is``
rather than parsing an e
Reading the task runtime context
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
-Declare an ``sdk.TIRunContext`` parameter on a task to read 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. It is an interface that embeds ``context.Context``, so the same
``ctx`` drives cancellation and
-client calls. The runtime binds it by type, just like the other injected
parameters:
+``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:
.. code-block:: go
- func extract(ctx sdk.TIRunContext, log *slog.Logger) (any, error) {
- ti := ctx.TaskInstance()
- log.Info("running",
+ 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", ctx.DagRun().LogicalDate,
+ "logical_date", actx.DagRun().LogicalDate,
)
return nil, nil
}
-``ctx.TaskInstance()`` returns ``DagID``, ``RunID``, ``TaskID``, ``MapIndex``
(nil for an unmapped task),
-and ``TryNumber``; ``ctx.DagRun()`` returns ``DagID``, ``RunID``, and the
``*time.Time`` fields
+``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).
+.. _go-sdk/arguments:
+
+Receiving arguments from the stub Dag
+~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
+
+Every parameter after the ``airflow.Context`` is a **data parameter**, filled
in declaration order from the
+arguments of the Python stub Dag'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.
+
+.. code-block:: go
+
+ // The stub Dag 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
+ }
+
+Stub parameters the Dag author left at their Python defaults are the
exception: 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.
Review Comment:
I'd like to move this part after mentioning the struct-base binding.
```suggestion
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -200,37 +204,59 @@ There is no separate Go worker to run: the Airflow worker
forks the bundle binar
Writing tasks
-------------
-The runtime inspects a task function's signature and injects arguments by
type, so you only declare the
-parameters your task actually needs:
+Every task function takes an ``airflow.Context`` as its first parameter, and
reaches what Airflow provides
+through its methods:
.. list-table::
:header-rows: 1
:widths: 35 65
- * - Parameter type
- - Injected value
- * - ``sdk.TIRunContext``
- - The task's execution context: the cancellation/deadline signal plus the
task instance identifiers and
- Dag run timestamps. Respect it for long-running work. See
:ref:`go-sdk/runtime-context`.
- * - ``*slog.Logger``
- - 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.
+ * - Method
+ - What it returns
+ * - ``actx.Logger()``
+ - An ``*slog.Logger`` whose output is routed back to the Airflow task log.
+ * - ``actx.Client()``
+ - A client for Airflow Variables, Connections, and XCom. See
:ref:`go-sdk/client`.
+ * - ``actx.TaskInstance()``
+ - The identifiers of the running task instance. See
:ref:`go-sdk/runtime-context`.
+ * - ``actx.DagRun()``
+ - The identifiers and scheduling timestamps of its Dag run. See
:ref:`go-sdk/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 Dag's TaskFlow
call; see
Review Comment:
```suggestion
Every parameter after the Context is data, filled from the stub Task's
TaskFlow call; see
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst:
##########
@@ -268,30 +295,74 @@ so you can branch on a missing value with ``errors.Is``
rather than parsing an e
Reading the task runtime context
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
-Declare an ``sdk.TIRunContext`` parameter on a task to read 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. It is an interface that embeds ``context.Context``, so the same
``ctx`` drives cancellation and
-client calls. The runtime binds it by type, just like the other injected
parameters:
+``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:
.. code-block:: go
- func extract(ctx sdk.TIRunContext, log *slog.Logger) (any, error) {
- ti := ctx.TaskInstance()
- log.Info("running",
+ 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", ctx.DagRun().LogicalDate,
+ "logical_date", actx.DagRun().LogicalDate,
)
return nil, nil
}
-``ctx.TaskInstance()`` returns ``DagID``, ``RunID``, ``TaskID``, ``MapIndex``
(nil for an unmapped task),
-and ``TryNumber``; ``ctx.DagRun()`` returns ``DagID``, ``RunID``, and the
``*time.Time`` fields
+``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).
+.. _go-sdk/arguments:
+
+Receiving arguments from the stub Dag
+~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
+
+Every parameter after the ``airflow.Context`` is a **data parameter**, filled
in declaration order from the
+arguments of the Python stub Dag'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.
+
Review Comment:
How about adding something like (please feel free to rephrase my statement)
```suggestion
There're two ways for receiving the arguments from Python Stub Task side:
1. positional binding
2. struct-based (keyword-based) binding
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]