TypeScript SDK - Dag/DagRegistry 구조와 빌드·배포
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/language-sdks/typescript.rst - TypeScript SDK (전체, 실험적/베타, 약 324줄)
이 모듈을 다 읽으면
- TypeScript SDK에서 Dag/DagRegistry/serveDags가 각각 어떤 역할을 하는지 설명할 수 있다
- TaskClient의 XCom/Variable/Connection 조회 API와 null 처리 방식을 설명할 수 있다
- airflow-ts-pack이 만드는 bundle.mjs의 메타데이터 내장 방식과 esbuild의 위치를 설명할 수 있다
TypeScript SDK(apache-airflow-ts-sdk, ESM 전용, 베타)는 Node.js 22 위에서 동작하며, Dag에 태스크 핸들러를 등록하고 DagRegistry로 모아 serveDags로 서빙하는 구조이고, airflow-ts-pack이 esbuild로 엔트리 모듈과 의존성을 하나의 bundle.mjs로 묶으면서 매니페스트를 base64 주석으로 embed한다.
Dag / DagRegistry / serveDags 구조
TypeScript SDK는 `apache-airflow-ts-sdk` 패키지(ESM 전용)로 배포되며 현재 베타 상태라 API가 바뀔 수 있다. 태스크 핸들러는 `TaskHandlerArgs`를 받는 평범한(보통 `async`) 함수다. 그 `dag_id`를 구현하는 `Dag`를 만들고, `dag.task`로 각 핸들러를 붙인 뒤, 여러 `Dag`를 `DagRegistry`로 모아 `serveDags`로 서빙한다 - 최상위 `await`로 쓰이는 이 호출이 모듈 자체를 실행 가능한 번들 진입점으로 만든다.
`new Dag(...)`에 넘기는 `dagId`는 파이썬 Dag의 `dag_id`와 일치해야 하고, `dag.task`에 넘기는 `taskId`는 그 Dag의 `@task.stub` 함수 이름과 일치해야 한다. `serveDags`에 넘기는 레지스트리가 그 번들이 담는 완전한 Dag 집합이다 - `serveDags`를 두 번째로 호출하면 거부되고, 레지스트리에서 빠진 Dag는 패킹된 번들에 포함되지 않아 그 태스크들은 런타임에 제거된 것으로 표시된다.
`DagRegistry`는 소켓을 열지도, 아무것도 시작하지도 않으므로, 코디네이터 런타임 없이 유닛 테스트에서 `registry.getTaskHandler(dagId, taskId)`로 핸들러를 직접 가져와 호출할 수 있다. 여러 모듈에 걸쳐 Dag를 모으는 번들이라면 `registry.register(...)`로 점진적으로 추가할 수 있다. `new Dag`와 `dag.task`는 후행 옵션 객체(`Dag`는 `spec`, 태스크는 `spec`과 `inputs`)를 받지만 아직 쓰이지 않으므로 설정해서는 안 되고, 그 외 키를 넘기면 거부된다.
핵심 포인트
- Dag에 dag.task(taskId, handler)로 핸들러를 붙이고 여러 Dag를 DagRegistry로 모아 serveDags로 서빙하는데, serveDags는 한 번만 호출할 수 있고 레지스트리에 빠진 Dag는 번들에 포함되지 않아 그 태스크들이 런타임에 제거된 것으로 표시된다
- DagRegistry는 소켓을 열거나 아무것도 시작하지 않으므로, 코디네이터 런타임 없이 registry.getTaskHandler(dagId, taskId)로 유닛 테스트에서 핸들러를 직접 호출할 수 있다
- new Dag/dag.task가 받는 spec/inputs 같은 후행 옵션 객체는 아직 쓰이지 않으므로 설정하지 말아야 하고, 그 외 키를 넘기면 거부된다
TaskHandlerArgs와 TaskClient
모든 태스크 핸들러는 `TaskHandlerArgs` 객체 하나를 받는다. `ctx` 필드는 `dagId`, `taskId`(TaskGroup 접두사 포함), `runId`, `tryNumber`, `mapIndex`(매핑되지 않은 태스크는 `-1`), 그리고 Airflow가 태스크를 종료할 때 발동하는 `AbortSignal`인 `signal`을 담는다 - `signal`을 `fetch()`나 타이머 등 `AbortSignal`을 받는 API에 넘기면 협조적 취소(cooperative cancellation)를 구현할 수 있다. `client` 필드는 Variable/Connection/XCom을 위한 `TaskClient`다.
핸들러가 `undefined`가 아닌 값을 반환하면 파이썬 `@task`와 동일하게 그 값이 `return_value` XCom이 되고, 잡히지 않은 예외나 reject된 프라미스는 태스크 인스턴스를 실패로 표시해 stub에 재시도가 설정되어 있으면 트리거한다.
`TaskClient`의 표면은 다음과 같다 - `getVariable(key)`는 값이 없으면 `null`을 반환하는 관대한 버전이고, `getVariableOrThrow(key)`는 대신 `VariableNotFoundError`를 던져 기본값 없는 파이썬 `Variable.get`과 동작을 맞춘다. `getConnection(connId)`는 `id`, `type`과 선택 필드 `host`, `schema`, `login`, `password`, `port`, `extra`(각각 없거나 `null`일 수 있음)를 가진 `ConnectionResult`를 반환하거나, 연결이 없으면 `null`을 반환하며, `getConnectionOrThrow(connId)`는 대신 `ConnectionNotFoundError`를 던져 파이썬 `BaseHook.get_connection`과 동작을 맞춘다. `getXCom<T>({key, ...})`는 XCom 값을 읽거나 없으면 `null`을 반환하며, 로케이터 필드(`dagId`, `runId`, `taskId`, `mapIndex`)는 기본적으로 현재 태스크를 가리키고 `taskId`를 넘기면 업스트림 태스크의 XCom을 읽을 수 있다. `setXCom({key, value, ...})`은 XCom 값을 발행한다.
핵심 포인트
- ctx.signal은 Airflow가 태스크를 종료할 때 발동하는 AbortSignal이므로 fetch()나 타이머에 넘겨 협조적 취소(cooperative cancellation)를 구현할 수 있다
- getVariable/getConnection/getXCom은 값이 없으면 null을 반환하는 관대한 버전이고, getVariableOrThrow/getConnectionOrThrow는 각각 Python의 Variable.get(기본값 없음)/BaseHook.get_connection과 동작을 맞춰 예외를 던진다
- 태스크 핸들러가 undefined가 아닌 값을 반환하면 파이썬 @task와 동일하게 그 값이 return_value XCom이 되고, 처리되지 않은 예외나 reject된 프라미스는 태스크를 실패로 표시해 재시도를 트리거한다
XCom 타입 매핑, 빌드(airflow-ts-pack)와 코디네이터 설정
XCom은 JSON으로 저장되며 Python `int`/`float` 모두 JavaScript에서는 `number`로 매핑된다 - JavaScript는 숫자 타입이 IEEE 754 double인 `number` 하나뿐이라 정수와 소수가 같은 타입으로 오며, `Number.MAX_SAFE_INTEGER`(2^53 − 1)를 넘는 정수는 정밀도를 잃을 수 있다.
로깅은 태스크가 `console.log`/`console.error`로 stdout/stderr에 쓰는 내용이 워커에 의해 캡처되어 Airflow 태스크 로그에 표시된다(stdout은 INFO 레벨, stderr는 ERROR 레벨) - 아직 전용 구조화 로깅 API는 없다.
빌드는 SDK와 함께 배포되는 `airflow-ts-pack`이 담당한다 - esbuild로 엔트리 모듈과 그 모든 import를 하나의 자기 완결적 ESM 파일 `bundle.mjs`로 묶고, 매니페스트(`dag_id`/`task_id` 맵과 슈퍼바이저 스키마 버전)를 선행 `//# airflowMetadata=<base64>` 주석으로 embed한다 - 별도의 매니페스트나 `node_modules` 없이 배포할 파일이 하나로 끝난다. `esbuild`는 선택적 peer 의존성이라 `apache-airflow-ts-sdk`의 런타임 설치에는 포함되지 않고, `airflow-ts-pack` 실행 전에 별도로 설치해야 한다(`npm install --save-dev esbuild`). `--outdir <dir>`로 출력 디렉터리(기본 `dist`)를, `--source <name>`으로 Airflow UI에 표시할 소스 이름(기본은 엔트리 파일의 basename)을 지정할 수 있다.
배포는 `bundle.mjs`를 코디네이터의 `bundles_root`로 지정된 디렉터리에 복사하거나 마운트하면 된다. `NodeCoordinator`는 설정된 디렉터리들을 순서대로 검색해 처음 발견한 사용 가능한 번들을 `node`로 실행한다. `NodeCoordinator`의 kwargs는 `bundles_root`(필수, embed된 메타데이터를 가진 `bundle.mjs`를 순서대로 검색할 하나 이상의 디렉터리), `node_executable`(기본값 `"node"`, PATH에서 찾음), `task_startup_timeout`(기본값 10.0초)이다.
알려진 제약으로는, 파이썬 stub Dag가 여전히 필수이고(실행 API가 아직 비-파이썬 언어의 Dag 구조를 나르지 못함), API가 호환되지 않게 바뀔 수 있는 베타 상태이며, `NodeCoordinator`는 `bundles_root`에서 찾은 첫 번째 사용 가능한 번들 하나만 실행한다(다른 Dag/태스크를 다른 번들로 라우팅하지 않으므로 여러 번들을 서빙하려면 별도 큐마다 코디네이터를 여러 개 등록해야 함). 또한 태스크 인스턴스마다 Node.js 서브프로세스가 하나씩 뜨므로 인스턴스 간 공유 상태가 필요한 태스크는 XCom이나 외부 저장소를 써야 한다.
핵심 포인트
- JavaScript는 숫자 타입이 number 하나뿐이라 정수와 소수가 같은 타입으로 오며, Number.MAX_SAFE_INTEGER(2^53-1)를 넘는 정수는 정밀도를 잃을 수 있다
- airflow-ts-pack은 esbuild로 엔트리 모듈과 의존성을 bundle.mjs 하나로 묶고 매니페스트를 //# airflowMetadata=<base64> 주석으로 embed한다 - esbuild는 선택적 peer 의존성이라 런타임 설치에는 포함되지 않고 패킹 전에 별도로 설치해야 한다
- NodeCoordinator는 bundles_root에서 순서대로 첫 번째로 쓸 수 있는 번들 하나만 실행한다 - 여러 번들을 서빙하려면 별도 큐마다 코디네이터를 여러 개 등록해야 한다