TypeScript SDK
This is an experimental feature.
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:
Dag definition. The schedule, the tasks, their options and the dependencies between them are all written in TypeScript, with no Python file involved. See Declaring a Dag in TypeScript.
Python Dag with stub TaskHandler. A Python Dag declares the tasks with
@task.stuband wires them together, and a TypeScriptTaskHandlerimplements each one. See Implementing stub tasks of a Python Dag.
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.
Warning
Install an available release from npm. To try an unreleased change, build it from source in the
ts-sdk/ directory of the Airflow repository
and depend on it locally (see ts-sdk/example/ for a working setup).
See also
For the full API reference (Dag, Bundle, TaskHandler, withArgNames, the task getters,
TaskClient, supporting types, and exceptions), see the
TypeScript SDK API reference.
Prerequisites
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. To add Node.js to the Airflow image, extend it as described in Building the image. If
nodeis not on thePATH, set the coordinator’snode_executable(see NodeCoordinator configuration).The
apache-airflow-task-sdkpackage, installed with Airflow, providesNodeCoordinator, which runs the TypeScript code. No other Python package is needed.In your TypeScript project, install the SDK, and
esbuildto build the bundle:npm install apache-airflow-ts-sdk npm install --save-dev esbuild
Quick start
This example declares a whole Dag in TypeScript, builds it into a bundle, and deploys it, with no Python file involved.
Declare a Dag in src/main.ts:
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:
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 Dag bundles).
Then send the typescript queue to the Node.js coordinator in airflow.cfg:
[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.
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:
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:
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
XCom type mapping).
A few rules keep the graph honest:
Every task is called exactly once. A task that is never called fails when the Dag is read, so none is left out of the graph by accident.
A reference has to be the input itself. One placed inside an array or an object is treated as a literal and draws no dependency.
Task ids
The task id may be left out, in which case it is the handler’s function name:
const extract = dag.task(async function extract(): Promise<number> {
return 42;
});
airflow-ts-pack keeps function names, so building the bundle never renames a task. A handler with no name
of its own, such as an arrow function written inline, needs an id: pass it first, or set taskId in the
task’s options. Give it in one place only.
A task id you write is made of letters, digits, dashes and underscores.
Dag and task options
new Dag and dag.task both take a trailing object of Airflow options. The fields are Airflow’s own,
spelled in camelCase, and a misspelled one is a compile error:
const dag = new Dag("ts_etl", {
schedule: "@daily",
catchup: false,
tags: ["etl"],
queue: "typescript",
});
const load = dag.task("load", loadRows, { retries: 2, retryDelay: 30 });
schedule accepts @once, @continuous, a cron expression, or a cron preset such as @daily. Leave
it out for a Dag that only runs when triggered. A cron schedule runs in UTC.
Every task of the Dag runs on the Node.js coordinator, so it needs a queue that
queue_to_coordinator sends there. Set queue once on the Dag and
each task inherits it. A task’s own queue wins over the Dag’s.
Order-only dependencies
When a task only has to run after another, without receiving its value, draw the dependency between the
references with before and after, the TypeScript spelling of Python’s >> and <<:
const loaded = load({ total });
const cleaned = cleanup();
loaded.before(cleaned); // loaded >> cleaned
cleaned.after(loaded, extracted); // [loaded, extracted] >> cleaned
Both take any number of references, and drawing a dependency that already exists changes nothing. Each returns
the reference it was called on, so loaded.before(cleaned).before(notified) draws both from loaded.
Pass a value as an input when the downstream task needs it, and use before or after when it only needs
to run in order.
Task groups
dag.taskGroup(groupId) opens a group with the same task and taskGroup methods as the Dag. Every id
declared in it is prefixed with the group’s, as in Python:
const staging = dag.taskGroup("staging");
const staged = staging.task("stage_rows", stageRows)(); // task id "staging.stage_rows"
staging.taskGroup("checks").task("nulls", checkNulls)(); // task id "staging.checks.nulls"
staging.before(loaded); // staging >> loaded
A group takes before and after too, which orders the whole group against a task or another group.
Tasks and groups share one id namespace, so a Dag cannot hold both a task and a group called staging. The
. is added by the group, so an id you write cannot contain one.
Pass { prefixGroupId: false } to keep the ids declared in a group as written, as prefix_group_id=False
does in Python; they then have to be unique across the Dag. A group id is made of letters, digits, dashes and
underscores, and is at most 200 characters.
Conditional branching
dag.if declares a task from a handler that returns a boolean, and names the task each outcome runs:
async function hasRows({ total }: { total: number }): Promise<boolean> {
return total > 0;
}
dag.if(hasRows, { total }).then(loaded).else(reportedEmpty);
Its inputs come second, and an optional third argument takes task options, such as taskId; the id is
otherwise the handler’s name. The compiler checks that the handler returns a boolean. else is optional:
without it, the then task is skipped when the condition is false.
The tasks a condition guards take no input from it. To use a value the condition computed, read it with
getClient().getXCom.
The side not taken is skipped, and stays skipped if it is cleared later. Only the tasks named here are
skipped, so a task that both sides lead to still runs. This differs from Python’s @task.branch, which
skips every direct downstream task it did not choose.
Multi-way branching
dag.switch declares a task from a handler that returns the reference of the task to run, and lists the
candidates:
const daily = publishDaily();
const weekly = publishWeekly();
async function pickCadence() {
const cadence = await getClient().getVariable("cadence");
return cadence === "weekly" ? weekly : daily;
}
dag.switch(pickCadence).case(daily).case(weekly);
It takes inputs and options as dag.if does. A case is a task reference, so the compiler checks it exists,
and renaming a handler cannot silently change which task runs. The deciding task’s value is the id of the
task it chose.
Exactly one case runs, and there is no default: a decider that returns anything else fails, naming what it
returned and what it could have returned. To run several tasks together on one outcome, put them behind a
single task, or give each its own dag.if.
Triggering another Dag
triggerDagRun makes a task that starts a run of another Dag. Pass it to dag.task in place of a handler:
import { triggerDagRun } from "apache-airflow-ts-sdk";
const triggered = dag.task(
"trigger_downstream",
triggerDagRun({ dagId: "downstream_etl", conf: { source: "ts_etl" }, waitForCompletion: true }),
)();
triggered.after(loaded);
The task has no handler to take an id from, so it needs one, and it takes no inputs. triggerDagRun takes
the options of Python’s TriggerDagRunOperator, such as conf, waitForCompletion, deferrable,
pokeInterval, allowedStates and failedStates. Values are sent as written, so Jinja in conf is
not rendered.
The task runs on Node.js like the Dag’s other tasks, and inherits the Dag’s queue. With deferrable, it
waits in the triggerer, which needs the standard provider installed.
A complete example
ts-sdk/example/src/native.ts uses all of the above in one Dag: a task group, a fan-in, a condition, a
multi-way branch, order-only dependencies and a triggered Dag run.
Implementing stub tasks of a Python Dag
Here a Python Dag declares the tasks and their dependencies, and TypeScript implements some of them. Use it when the Dag needs something a TypeScript Dag cannot express yet (see Limitations), or to mix Python and TypeScript tasks in one Dag.
The Python Dag
from airflow.sdk import dag, task
@dag
def typescript_example():
@task
def python_start():
return "hello from Python"
@task.stub(queue="typescript")
def build_message(): ...
python_start() >> build_message()
typescript_example()
@task.stub declares a task that is implemented elsewhere, and queue sends it to the Node.js
coordinator.
The TypeScript side
Create a TaskHandler per stub task, naming the dag_id and task_id it implements, and register them
on a Bundle:
import { Bundle, getClient, TaskHandler } from "apache-airflow-ts-sdk";
export async function buildMessage() {
const client = getClient();
const upstream = await client.getXCom<string>({ key: "return_value", taskId: "python_start" });
const greeting = await client.getVariable("typescript_example_greeting");
return `${greeting ?? "hello from TypeScript"}; upstream=${upstream ?? "missing"}`;
}
const bundle = new Bundle();
bundle.register(new TaskHandler("typescript_example", "build_message", buildMessage));
await bundle.serve();
The taskId must match the stub’s task id, including any task group prefix.
register takes any number of task handlers and Dags, so one bundle can implement the tasks of several
Python Dags and declare Dags of its own. A task left out of the bundle is marked removed when it runs.
Nothing connects to Airflow until bundle.serve(), so a unit test can build a bundle and call a handler
through bundle.getTaskHandler(dagId, taskId) without running Airflow.
TaskFlow arguments
A Python Dag that calls a stub task TaskFlow-style passes those arguments to the handler, which destructures them by name:
@task.stub(queue="typescript")
def transform(region_code: str, threshold: float, dry_run: bool = False): ...
transform("uk", 0.75)
interface TransformArgs {
regionCode: string;
threshold: number;
dryRun: boolean;
}
export async function transform({ regionCode, threshold, dryRun }: TransformArgs) {
// ...
}
Names match ignoring case and underscores, so region_code reaches regionCode and s3_uri reaches
s3Uri with nothing declared on either side. The Go SDK matches names the same way. An argument the call
leaves at its default arrives with the default’s value.
A name the handler reads that the call did not pass is logged rather than raised, naming what the handler asked for and what the call delivered. Two Python names that match the same TypeScript name fail the task.
An argument filled from another task, as in transform(extract(), "uk"), arrives as that task’s value. An
upstream task that produced no value fails the task, naming the argument and the task; one that returned
None delivers null.
A Python int larger than a JavaScript number holds exactly (beyond ±9007199254740991) is refused, so pass
such a value as a string.
Renaming an argument
withArgNames maps a handler’s name to a different Python name, for when the two cannot match on their own:
a clearer word than the Dag chose, or a TypeScript reserved word such as enum. The map comes first, the
handler second:
interface ReportArgs {
summary: Summary;
label: string; // the Python call passes this as `run_label`
}
const report = withArgNames({ label: "run_label" }, async ({ summary, label }: ReportArgs) => {
// ...
});
bundle.register(new TaskHandler("etl", "report", report));
Every name the map does not mention still matches as usual, so withArgNames should be rare. The map’s keys
are checked against the handler’s argument type, so { labl: "run_label" } is a compile error.
Note
A dependency drawn with >> in the Python Dag orders the tasks and passes nothing. To read a value the
call did not pass, use getClient().getXCom.
Writing tasks
A handler is an ordinary, usually async, function. It works the same in both ways of using the SDK.
Two functions give it what it needs while it runs:
Function |
Returns |
|---|---|
|
The task’s context (a |
|
A |
Both work anywhere inside the handler’s call, including helper functions and code after an await, and
throw outside it. Await everything the handler starts: work still running after the handler returns runs
after Airflow has recorded the task’s result.
A value the handler returns becomes the task’s return_value XCom, as with Python’s @task, and is what a
downstream task receives. An uncaught exception or a rejected promise fails the task, and retries it if the task
has retries.
The TaskClient
getVariable(key)returns the Variable as a string, ornullwhen it does not exist.getVariableOrThrow(key)throwsVariableNotFoundErrorinstead, like Python’sVariable.getwithout a default.setVariable(key, value, description?)stores a Variable, replacing any existing value, anddeleteVariable(key)removes one. Values are strings, so serialize structured data, for example withJSON.stringify, before storing it.getConnection(connId)returns aConnectionResultwithidandtype, plushost,schema,login,password,portandextra, each of which may be missing ornull. It returnsnullwhen the connection does not exist, andgetConnectionOrThrow(connId)throwsConnectionNotFoundErrorinstead, like Python’sBaseHook.get_connection.getXCom<T>({ key, ... })reads an XCom value, ornullwhen it does not exist.dagId,runId,taskIdandmapIndexdefault to the current task; passtaskIdto read an upstream task’s XCom. See XCom type mapping for how values map to JavaScript types.setXCom({ key, value, ... })publishes an XCom value.taskStateStoreis this task instance’s key/value store, withget<T>(key),set(key, value, { retentionMs? }),delete(key)andclear().
Note
A value from a secrets backend, such as an AIRFLOW_VAR_* environment variable, still wins over the
stored value when the Variable is read back. setVariable without a description clears the existing
description, and deleteVariable succeeds even when the Variable does not exist.
Task state store
The task state store is a per-task-instance key/value store, scoped to dagId, runId, taskId,
and mapIndex (but not tryNumber), so a value written by one attempt is still readable by the next.
It survives worker crashes and task retries within the same Dag run, which makes it a good place to record
external job IDs, checkpoints within a task, and progress metadata. See Task State Store
for the full concept, including the equivalent Python API.
import { getClient, NEVER_EXPIRE } from "apache-airflow-ts-sdk";
export async function runSparkJob() {
const store = getClient().taskStateStore;
let jobId = await store.get<string>("job_id");
if (jobId == null) {
jobId = await submitSparkJob();
await store.set("job_id", jobId, { retentionMs: NEVER_EXPIRE });
}
const result = await waitForSparkJob(jobId);
await store.delete("job_id");
return result;
}
Retention. retentionMs is milliseconds to retain the key, counted from the time of the write:
A number retains the key for that many milliseconds from now;
retentionMs: 0expires the key immediately.NEVER_EXPIRE(imported fromapache-airflow-ts-sdk) stores the key with no expiry, regardless of the deployment default.When
retentionMsis omitted the runtime readsAIRFLOW__STATE_STORE__DEFAULT_RETENTION_DAYS, which the coordinator sets from[state_store] default_retention_days. A value of0there means the key never expires. This is not the same asretentionMs: 0above, which expires the key immediately.
Keys must be non-empty strings of at most 512 characters; they may contain slashes. Values must be
JSON-compatible (see XCom type mapping), and can never be null — set rejects a null
value. get returns null only to mean “key not found.”
Two behaviors to note:
[workers] state_store_backendis not used from TypeScript tasks, so values are always stored inline through the Execution API. A key written by a Python task through a custom worker-side backend reads back from TypeScript as that backend’s reference marker string, not the original value.[state_store] clear_on_successremoves all of a task instance’s state store keys when the task succeeds, the same as for Python tasks.
Logging
Anything a task writes to stdout or stderr, such as console.log or console.error, appears in the
Airflow task log: stdout at INFO level, stderr at ERROR. There is no dedicated structured-logging API
yet.
XCom type mapping
XComs are stored as JSON, so a value read with getXCom arrives as the matching JavaScript type:
Python type |
JSON |
JavaScript type |
|---|---|---|
|
number (integer) |
|
|
number (decimal) |
|
|
string |
|
|
boolean |
|
|
null |
|
|
array |
|
|
object |
|
Note
JavaScript has a single number type, so integers and decimals arrive as the same type, and an integer
larger than Number.MAX_SAFE_INTEGER (253 − 1) may lose precision.
Building a bundle
airflow-ts-pack, shipped with the SDK, builds your entry module and everything it imports into a single
file, bundle.min.mjs. That file is all you deploy: it needs no node_modules and no separate metadata
file, and Airflow refuses to run a bundle whose content was changed after it was built.
npx airflow-ts-pack src/main.ts --outdir dist
--outdir <dir>sets the output directory (defaultdist).--outfile <path>names the output file exactly, which helps when one directory holds several bundles. The name must end in.min.mjs.
--outdir and --outfile cannot be used together. esbuild must be installed before you build; the
runtime install of apache-airflow-ts-sdk does not need it.
The bundle’s code is minified, with license banners of bundled dependencies kept. The Airflow UI’s Code view shows the file each Dag is declared in, as you wrote it.
Deploying
Where the bundle goes depends on what it contains.
Dag definition
Put the bundle in a Dag bundle, such as the default dags-folder bundle. The Dag processor runs each
*.min.mjs bundle it finds with node to read its Dags, as it reads any other Dag file, so the Dag processor
needs Node.js and the same [sdk] configuration as the workers:
[sdk]
coordinators = {"ts": {"classpath": "airflow.sdk.coordinators.node.NodeCoordinator"}}
queue_to_coordinator = {"typescript": "ts"}
With more than one NodeCoordinator, map each Dag bundle that holds *.min.mjs bundles to one of them in
[sdk] dag_bundle_to_coordinator, such as {"dags-folder": "ts"}, or its bundles fail to parse.
Their tasks run from the same Dag bundle, at the version their Dag run was created with, so a Dag declared in TypeScript needs no coordinator option to find its bundle.
A bundle that fails its integrity check, or whose Dags cannot be read or have a cycle, shows up as an import
error, and the other Dags keep working. Do not declare the same dag_id in a Python file of that Dag bundle
too: the two overwrite each other on every parse.
Python Dag with stub TaskHandler
By default, the coordinator loads the bundle that implements a stub task from the stub task’s own Dag bundle, at the version its Dag run was created with, so the bundle goes next to the Python Dag that declares the stub.
To keep such bundles in a Dag bundle of their own, name it with task_handler_bundle_name. A task then uses
the version that bundle is on when the task starts:
[sdk]
coordinators = {
"ts": {
"classpath": "airflow.sdk.coordinators.node.NodeCoordinator",
"kwargs": {"task_handler_bundle_name": "ts-handlers"}
}
}
queue_to_coordinator = {"typescript": "ts"}
File names do not matter beyond the .min.mjs suffix, so one Dag bundle can hold several bundles.
The Dag processor also reads a bundle that only implements stub tasks, which starts node once per parse
and finds no Dags. List such bundles in .airflowignore to skip that; the coordinator still finds them.
Where the configuration is needed
The coordinator runs where tasks run, so the [sdk] configuration and Node.js are needed on the workers.
With CeleryExecutor, set them on the Celery workers. With LocalExecutor, tasks run in the scheduler’s
process, so set them there. The Dag processor needs them to read Dags declared in TypeScript. If it has the
[sdk] configuration, it needs Node.js even without such Dags, since it runs every *.min.mjs bundle it
finds. The API server never needs them.
There is no separate Node.js service to run: the worker starts the bundle with node once per task
instance.
NodeCoordinator configuration
The kwargs of a coordinators entry are passed to
NodeCoordinator, and all of them are optional:
Parameter |
Default |
Description |
|---|---|---|
|
(unset) |
The Dag bundle to load the bundles that implement stub tasks from. When unset, a stub task loads its
bundle from its own Dag bundle. It must be registered in |
|
|
Path to the |
|
|
Seconds to wait for a task’s bundle to start. Increase it if your bundle starts slowly, for example on constrained hardware. |
task_handler_bundle_name only concerns stub tasks: a Dag declared in TypeScript always runs from
the Dag bundle it was read from.
queue_to_coordinator maps a queue to a coordinator entry, so every task on that queue runs on Node.js.
Limitations
A Dag declared in TypeScript cannot express everything yet. Assets, custom timetables, dynamic task mapping and setup and teardown tasks have no TypeScript form. Declare such a Dag in Python, and implement its tasks in TypeScript with
@task.stub.Cluster policies do not apply to Dags declared in TypeScript.
dag_policyandtask_policydo not run when the Dag processor reads them.Some CLI commands do not take a Dag declared in TypeScript.
airflow dags test,tasks testandtasks renderrefuse it, andairflow dags reserializeskips it.Beta. The API may change in incompatible ways between releases.
One Node.js process per task instance. Tasks cannot share in-process state; use XComs or an external store.