pierrejeambrun commented on code in PR #73875:
URL: https://github.com/apache/airflow/pull/73875#discussion_r4131749989
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -22,12 +22,17 @@ TypeScript SDK
|experimental|
-The TypeScript SDK lets you register task handlers on a ``Bundle`` and
implement their logic in TypeScript (or
-plain JavaScript), running on Node.js. A matching Python stub Dag still
declares the scheduling shape and
-dependencies; individual tasks delegate to a Node.js subprocess that is
spawned by
-:class:`~airflow.sdk.coordinators.node.NodeCoordinator` for each task instance.
+The TypeScript SDK lets you write Airflow Dags and tasks in TypeScript, or
plain JavaScript, and run them on
+Node.js. There are two ways to use it:
-The SDK is the ``apache-airflow-ts-sdk`` package (ESM-only). It is currently
in **beta** and its API may change.
+* **Declare the whole Dag in TypeScript.** The schedule, the tasks, their
options and the dependencies between
+ them are all written in TypeScript, with no Python file involved. See
:ref:`typescript-sdk/native-dag`.
+* **Implement mixed-language tasks.** A Python Dag declares the tasks with
``@task.stub`` and wires them
+ together, and TypeScript supplies what each task does. See
:ref:`typescript-sdk/mixed`.
+
+Both use the same task API and the same build tool, and one bundle can serve
both.
+
+The SDK is the ``apache-airflow-ts-sdk`` npm package (ESM-only). It is in
**beta**, and its API may change.
Review Comment:
A detail but I think original targeted namespace for pnpm was recently
released (cf mailing list).
`apache-airflow-ts-sdk` is good too if we want to continue that way.
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -48,25 +53,311 @@ The SDK is the ``apache-airflow-ts-sdk`` package
(ESM-only). It is currently in
Prerequisites
-------------
-* Node.js 22 or later must be available on the Airflow worker nodes.
-* The packed bundle (a single ``bundle.min.mjs`` file, see
:ref:`typescript-sdk/build`) must be accessible
- from the worker, under a directory the coordinator scans.
-* The ``apache-airflow-task-sdk`` package (installed with Airflow) provides
the coordinator; no additional
- Python packages are needed.
-* In the TypeScript project, install the ``apache-airflow-ts-sdk`` npm package
to author task handlers:
+* Node.js 22 or later on the Airflow workers. A Dag declared in TypeScript
also needs Node.js on the Dag
+ processor, which runs the bundle to read the Dag.
+* The ``apache-airflow-task-sdk`` package, installed with Airflow, provides
+ :class:`~airflow.sdk.coordinators.node.NodeCoordinator`, which runs the
TypeScript code. No other Python
+ package is needed.
+* In your TypeScript project, install the SDK, and ``esbuild`` to build the
bundle:
.. code-block:: bash
npm install apache-airflow-ts-sdk
+ npm install --save-dev esbuild
+
+.. _typescript-sdk/quick-start:
Quick start
-----------
-The following example shows the minimal moving parts: a Python Dag with a stub
task, and a TypeScript
-implementation of that task.
+Declare a Dag in ``src/main.ts``:
+
+.. code-block:: typescript
+
+ import { Bundle, Dag } from "apache-airflow-ts-sdk";
+
+ const dag = new Dag("ts_hello", { schedule: "@daily", queue: "typescript"
});
+
+ const extract = dag.task("extract", async (): Promise<number> => 42);
+ const report = dag.task("report", async ({ rows }: { rows: number }) => {
+ console.log(`extracted ${rows} rows`);
+ });
+
+ report({ rows: extract() });
+
+ await new Bundle(dag).serve();
+
+Build it into a single file:
+
+.. code-block:: bash
+
+ npx airflow-ts-pack src/main.ts --outdir dist
+
+Make sure ``dist/bundle.min.mjs`` is in a Dag bundle, for example the default
``dags-folder`` bundle, which
+reads the ``[core] dags_folder`` directory (see :doc:`Dag bundles
</administration-and-deployment/dag-bundles>`).
+Then send the ``typescript`` queue to the Node.js coordinator in
``airflow.cfg``:
+
+.. code-block:: ini
+
+ [sdk]
+ coordinators = {"ts": {"classpath":
"airflow.sdk.coordinators.node.NodeCoordinator"}}
+ queue_to_coordinator = {"typescript": "ts"}
+
+Airflow reads the bundle like any other Dag file: ``ts_hello`` shows up in the
UI, runs on its schedule,
+and ``report`` receives the value ``extract`` returned.
+
+.. _typescript-sdk/native-dag:
+
+Declaring a Dag in TypeScript
+-----------------------------
+
+Tasks and their inputs
+~~~~~~~~~~~~~~~~~~~~~~
+
+``dag.task(taskId, handler)`` declares a task and returns a *factory*. Calling
the factory places the task in
+the Dag and gives it its inputs, so the calls you write are the dependencies:
+
+.. code-block:: typescript
+
+ import { Dag } from "apache-airflow-ts-sdk";
+
+ const dag = new Dag("ts_etl");
+
+ const extract = dag.task("extract", async (): Promise<number> => 42);
+ const transform = dag.task(
+ "transform",
+ async ({ rows, region }: { rows: number; region: string }) => rows * 2,
+ );
+ const load = dag.task("load", async ({ total }: { total: number }) => {});
+
+ const extracted = extract();
+ const total = transform({ rows: extracted, region: "us" });
+ load({ total });
+
+A handler takes one object of named arguments, and a call names each input. An
input is either the
+reference another task's call returned, which makes this task wait for that
task and receive its value, or a
+literal JSON value such as ``"us"``. A task with no arguments is called with
none, and a single argument is
+named like any other: ``load({ total })``.
+
+The compiler checks every call: an argument left out, a misspelled one, and a
literal of the wrong type are
+all compile errors.
+
+The argument type can be a named interface, and the handler a function
declared anywhere, including in another
+module. This is also how a task passes data to the next one:
+
+.. code-block:: typescript
+
+ interface Summary {
+ total: number;
+ regions: number;
+ }
+
+ interface ReportArgs {
+ summary: Summary;
+ label: string;
+ }
+
+ export async function report({ summary, label }: ReportArgs) {
+ console.log(`${label}: ${summary.total} rows in ${summary.regions}
regions`);
+ }
+
+ const summarize = dag.task("summarize", async (): Promise<Summary> => ({
total: 42, regions: 2 }));
+ const reportTask = dag.task("report", report);
+
+ reportTask({ summary: summarize(), label: "nightly" });
+
+``report`` receives the object ``summarize`` returned as ``summary``, and the
literal ``"nightly"`` as
+``label``. Values move between tasks as XComs, so a task returns data JSON can
hold (see
+:ref:`typescript-sdk/types`).
+
+``withArgList`` gives the same inputs in order instead, for a call that reads
better that way:
Review Comment:
`withArgList` reference
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -37,9 +42,9 @@ The SDK is the ``apache-airflow-ts-sdk`` package (ESM-only).
It is currently in
.. seealso::
- For the full TypeScript API reference (``Bundle``, ``TaskHandler``, ``Dag``,
the task handler getters,
- ``TaskClient``, supporting types, and exceptions),
- see the `TypeScript SDK API reference
<https://airflow.apache.org/docs/ts-sdk/stable/>`__.
+ For the full API reference (``Dag``, ``Bundle``, ``TaskHandler``,
``withArgList``, ``withArgNames``, the
Review Comment:
Should remove the `withArgList` references added in this PR.
```suggestion
For the full API reference (``Dag``, ``Bundle``, ``TaskHandler``,
``withArgNames``, the
```
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -48,25 +53,311 @@ The SDK is the ``apache-airflow-ts-sdk`` package
(ESM-only). It is currently in
Prerequisites
-------------
-* Node.js 22 or later must be available on the Airflow worker nodes.
-* The packed bundle (a single ``bundle.min.mjs`` file, see
:ref:`typescript-sdk/build`) must be accessible
- from the worker, under a directory the coordinator scans.
-* The ``apache-airflow-task-sdk`` package (installed with Airflow) provides
the coordinator; no additional
- Python packages are needed.
-* In the TypeScript project, install the ``apache-airflow-ts-sdk`` npm package
to author task handlers:
+* Node.js 22 or later on the Airflow workers. A Dag declared in TypeScript
also needs Node.js on the Dag
+ processor, which runs the bundle to read the Dag.
+* The ``apache-airflow-task-sdk`` package, installed with Airflow, provides
+ :class:`~airflow.sdk.coordinators.node.NodeCoordinator`, which runs the
TypeScript code. No other Python
+ package is needed.
+* In your TypeScript project, install the SDK, and ``esbuild`` to build the
bundle:
.. code-block:: bash
npm install apache-airflow-ts-sdk
+ npm install --save-dev esbuild
+
+.. _typescript-sdk/quick-start:
Quick start
-----------
-The following example shows the minimal moving parts: a Python Dag with a stub
task, and a TypeScript
-implementation of that task.
+Declare a Dag in ``src/main.ts``:
Review Comment:
Keep the introduction sentence to explain what we are doing.
"implement a native dag in trypescript... "
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -48,25 +53,311 @@ The SDK is the ``apache-airflow-ts-sdk`` package
(ESM-only). It is currently in
Prerequisites
-------------
-* Node.js 22 or later must be available on the Airflow worker nodes.
-* The packed bundle (a single ``bundle.min.mjs`` file, see
:ref:`typescript-sdk/build`) must be accessible
- from the worker, under a directory the coordinator scans.
-* The ``apache-airflow-task-sdk`` package (installed with Airflow) provides
the coordinator; no additional
- Python packages are needed.
-* In the TypeScript project, install the ``apache-airflow-ts-sdk`` npm package
to author task handlers:
+* Node.js 22 or later on the Airflow workers. A Dag declared in TypeScript
also needs Node.js on the Dag
Review Comment:
`Node.js 22 or later on the Airflow workers` a doc reference or small hint
on how this can be achieved.
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -445,24 +681,33 @@ All ``kwargs`` in the ``coordinators`` config entry are
passed to the
* - Parameter
- Default
- Description
- * - ``bundles_root``
- - *(required)*
- - One or more directories searched recursively, in order, for an
integrity-verified ``*.min.mjs``
- bundle that declares the requested Dag. Accepts a string, a path, or a
list of strings/paths.
+ * - ``task_handler_bundle_name``
Review Comment:
Drops `bundles_root` for `task_handler_bundle_name`, where is that defined ?
(Future PR work?)
##########
ts-sdk/README.md:
##########
@@ -177,133 +126,56 @@ transform("uk", 0.75)
```
```ts
-interface TransformArgs {
- regionCode: string;
- threshold: number;
- dryRun: boolean;
-}
-
export async function transform({ regionCode, threshold, dryRun }:
TransformArgs) {
// ...
}
```
-Names bind by **folding on both sides**, lowercased with underscores removed,
so `region_code` reaches
-`regionCode` and `s3_uri` reaches `s3Uri` with nothing declared on either side.
-An argument the call leaves at its default arrives with the default's value.
-
-A name nothing folds to is **logged, not thrown**, naming what the handler
asked for and what the call
-delivered. Two Python names that fold to the same token fail the task.
-
-`Object.keys` and rest destructuring (`{ ...rest }`) yield Python's names, and
`in` folds like a read.
-
-### Upstream outputs
-
-An argument the Python call fills from another task arrives as that task's
value, not as a reference to it.
-Feeding the `transform` above from an upstream task adds one argument on each
side:
-
-```python
-# the Python Dag
-@task
-def extract() -> int: ...
-
-
[email protected](queue="typescript")
-def transform(rows: int, region_code: str, threshold: float, dry_run: bool =
False): ...
-
-
-transform(extract(), "uk", 0.75)
-```
-
-```ts
-interface TransformArgs {
- rows: number;
- regionCode: string;
- threshold: number;
- dryRun: boolean;
-}
-
-export async function transform({ rows, regionCode, threshold, dryRun }:
TransformArgs) {
- // `rows` is the number extract() returned.
-}
-```
-
-An upstream that pushed no output fails the task, naming both the argument and
the task it came from.
-An upstream that pushed `null` binds `null`.
+Names match ignoring case and underscores, so `region_code` reaches
`regionCode` with nothing declared on
+either side. An argument filled from another task, as in `transform(extract(),
"uk")`, arrives as that task's
+value. A dependency drawn with `>>` orders the tasks and passes nothing; read
such a value with
+`getClient().getXCom`.
-Being upstream is not the same as being passed. An XCom dependency declared
with `>>` defines task order only,
-so a value the call did not pass is read explicitly:
+`withArgNames` maps a handler's name to a different Python name, when the two
cannot match on their own:
```ts
-const rows = await getClient().getXCom<number>({ key: "return_value", taskId:
"extract" });
-```
-
-A Python `int` beyond the ±9007199254740991 a JavaScript number holds exactly
is refused rather than bound,
-so carry such a value across the boundary as a string.
-
-### Explicit renames
-
-`withArgNames` states a binding folding cannot reach, for a name the Python
side never used:
-a clearer word than the Dag chose, or a TypeScript reserved word like `enum`.
Mapping first, handler second:
-
-```ts
-interface ReportArgs {
- summary: Summary;
- label: string; // Python calls this `run_label`
-}
-
const report = withArgNames({ label: "run_label" }, async ({ summary, label }:
ReportArgs) => {
- // `label` is the call's `run_label`; `summary` folded as usual.
+ // `label` is the call's `run_label`.
});
-
-bundle.register(new TaskHandler("etl", "report", report));
```
-An entry beats folding, and everything the map does not mention still folds,
-so `withArgNames` should be rare in a real Dag.
-The map's keys are checked against the handler's own parameter type,
-so `{ labl: "run_label" }` is a compile error naming the right key.
-Its values are Python names, which `tsc` cannot see and does not check.
+## Writing tasks
-`Dag` is another interface, for a Dag declared natively in TypeScript, and is
still a work in progress.
+A handler is a plain, usually `async`, function. `getContext()` and
`getClient()` give it the task's context and
+Airflow access while it runs, so it takes no SDK-supplied argument. A value it
returns becomes the task's
+`return_value` XCom, and an uncaught error fails the task.
-Airflow launches the bundled entrypoint with `--comm=host:port` and
-`--logs=host:port`. `bundle.serve()` connects to those sockets, receives the
-task startup message, finds the registered handler for the Dag/task pair, and
-reports the terminal task state back to Airflow.
+## Building and deploying
-See [`example/`](https://github.com/apache/airflow/tree/main/ts-sdk/example)
for
-a coordinator-runtime example that packs a bundle with `airflow-ts-pack` and
-uses a Python stub Dag.
-
-## Packing bundles
-
-`airflow-ts-pack` produces a single self-contained bundle in one command.
-Packing is build-time only, so `esbuild` is an optional peer dependency the
-runtime install skips:
+`airflow-ts-pack` builds the entry module and everything it imports into one
file, `dist/bundle.min.mjs`:
```bash
-npm install --save-dev esbuild
-airflow-ts-pack src/main.ts --outdir dist
+npx airflow-ts-pack src/main.ts --outdir dist
```
-It bundles the entrypoint into a minified `dist/bundle.min.mjs` with esbuild,
then runs that bundle with
-`--airflow-metadata` so it reports its own registered Dag/task pairs and
supervisor schema version. The manifest is
-embedded as a compact JSON `//# airflowMetadata=...` comment after a leading
compact JSON `//# airflowBundle=...`
-layout descriptor, and the entry module is embedded verbatim in a `/*#
airflowSource ... #*/` block comment so
-Airflow can show the source a bundle was authored from, which its minified
code no longer is. The CLI records the
-integrity metadata for all three regions in that descriptor, so a coordinator
that is handed a bundle whose content
-was replaced fails loudly instead of running it. The result is one deployable
file with no hand-written metadata
-sidecar.
+- `--outdir <dir>`: output directory (default `dist`)
+- `--outfile <path>`: exact output path, whose name must end in `.min.mjs`
+- `--source <name>`: the source file name shown in the Airflow UI (default:
the entry file's name)
Review Comment:
this will need adjustment if https://github.com/apache/airflow/pull/73723
goes through
--
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]