Airflowを導入すると、DAG内のPythonへ処理を書き足したくなります。
外部APIからデータを取得し、整形後にDBへ保存する。小さな処理でも、アプリ本体がTypeScriptなら、ドメイン知識はPythonとTypeScriptへ分散します。
その結果、型、テスト、修正箇所が二重化します。実装差分も生まれやすくなります。
そこで、あるモノレポでは役割を次のように分離しました。
- Airflow: 実行時刻、再試行、監視を担当する
- TypeScript worker: 取得、変換、保存などの実処理を担当する
DAGにはコマンドだけを書く
DAGは、ワークスペース内のスクリプトを起動するだけです。
ただし、BashOperatorがコマンドを実行する場所はスケジューラではありません。executorから割り当てられたAirflow worker上で実行されます。
以下の例では、Airflow workerのイメージにNode.jsとpnpmを含め、モノレポを/opt/appへ配置しています。Airflow 2.4以降を想定した簡略例です。
from datetime import timedelta
import pendulum
from airflow import DAG
from airflow.operators.bash import BashOperator
default_args = {
'retries': 3,
'retry_delay': timedelta(minutes=5),
'retry_exponential_backoff': True,
'max_retry_delay': timedelta(minutes=30),
}
with DAG(
dag_id='external_data_sync',
schedule='0 * * * *',
start_date=pendulum.datetime(2026, 1, 1, tz='UTC'),
catchup=False,
default_args=default_args,
):
BashOperator(
task_id='run_sync_worker',
bash_command='pnpm --filter @example/sync-worker run:once',
cwd='/opt/app',
env={
'DATABASE_URL': '{{ conn.app_database.get_uri() }}',
'EXTERNAL_API_KEY': '{{ conn.external_api.password }}',
},
append_env=True,
skip_on_exit_code=None,
)
このDAGが知っている内容は、次の項目だけです。
- 毎時0分に実行する
- 失敗時は最大3回まで再試行する
- 待機時間には指数バックオフを適用する
- workerの配置場所と起動コマンドを指定する
- 接続情報をAirflow Connectionから渡す
接続情報の実値は、Airflow ConnectionまたはSecrets Backendで管理します。DAGへ直接書かず、workerのログにも出力しません。
append_env=Trueは、Airflow workerが持つPATHなどを維持しながら、指定した環境変数を追加するために設定しています。これがない場合、envの内容だけで実行環境が置き換わります。
また、BashOperatorは終了コード99を既定でskippedとして扱います。worker側で任意の非0コードを失敗にしたい場合は、skip_on_exit_code=Noneを明示します。
実処理はTypeScriptへ集約する
実際の処理は、通常のworkerパッケージとして実装します。
packages/sync-worker/
├── src/fetch.ts
├── src/transform.ts
├── src/persist.ts
└── package.json
workerは、Webアプリと型、DBアダプター、テスト環境を共有できます。Airflowを起動せず、次のコマンドだけでローカル実行できます。
pnpm --filter @example/sync-worker run:once
ここでのrun:onceは、1回分の同期を終えたら終了するコマンドです。常駐プロセス用の起動スクリプトとは分けています。
この分離により、デバッグと単体テストをAirflowから切り離せます。
この構成で得られたこと
ロジックを1言語に保てる
ドメインロジックはTypeScriptへ集約します。Python側へ型や変換処理を再実装する必要はありません。
DAGがほとんど変わらない
処理内容が変わっても、起動コマンドを維持できる限り、DAG修正は不要です。変更対象は主に実行間隔と再試行条件です。
Airflowなしで再現できる
障害調査では、最初にworkerを単独実行できます。DAG読み込みやスケジューラ、executorまで再現する必要はありません。
代わりに失うもの
1つのBashOperatorへ処理をまとめると、Airflow上ではworker全体が1タスクに見えます。処理単位ごとの再実行が必要な場合、この構成は不向きです。
そのため、worker側には次の対応が必要です。
- 処理件数と所要時間を構造化ログへ出す
- 途中失敗は握りつぶさず、終了時に非0コードを返す
- 再実行へ耐えられる冪等性を持たせる
- worker全体を再試行するコストを見積もる
Airflowの再試行では、取得、変換、保存を先頭から実行します。取得時間が長い処理では、再取得の負荷やAPIレート制限も考慮が必要です。
大規模なファンアウトが必要なら、Airflowのタスクグラフを使う方が自然です。処理単位ごとに依存関係を管理したい場合も同様です。
BashOperatorで足りなくなったら
Airflow workerとNode.js環境の共有が難しくなった場合、TypeScript workerをコンテナイメージとして配布します。そのうえで、DockerOperatorまたはKubernetesPodOperatorから起動できます。
DAGが指定する内容はイメージ、コマンド、実行条件だけです。ビジネスロジックをTypeScript側へ置く方針は変わりません。
ただし、コンテナレジストリ、provider設定、起動時間などの運用コストは増えます。小規模な構成では、実行環境を固定したBashOperatorの方が単純です。
まとめ
Airflowが得意なのは、いつ実行するかと失敗時にどう扱うかです。
一方、何を実行するかはアプリ側へ置きます。型、テスト、依存関係を共有できるため、管理しやすくなります。
DAGを薄い実行境界にすると、Airflowはオーケストレーションへ集中できます。ビジネスロジックもTypeScript側で素直に育てられます。
参考