henry3260 commented on code in PR #73528:
URL: https://github.com/apache/airflow/pull/73528#discussion_r4070972240
##########
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:
I put this at the start of the chapter intentionally placing it after the
details would lose its guiding effect for readers.
--
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]