Jobs プラグイン
Jobs プラグイン
AppKit アプリケーションから Databricks Lakeflow Jobs をトリガーし、監視します。
主な機能:
- 名前付きjobキーによる複数jobのサポート
- 環境変数からのjobの自動検出
- SSE ストリーミングでステータスを更新する run-and-wait
- Zod スキーマによるパラメータの検証
- タスクタイプに応じたパラメータのマッピング (notebook、python_wheel、sql など)
基本的な使い方
import { createApp, server, jobs } from "@databricks/appkit";
await createApp({
plugins: [server(), jobs()],
});明示的な jobs 設定がない場合、pluginは環境変数 DATABRICKS_JOB_ID を読み取り、default キーとして登録します。
設定オプション
| オプション | 型 | デフォルト | 説明 |
|---|---|---|---|
timeout | number | 60000 | Jobs API 呼び出しのデフォルトタイムアウト (ミリ秒) |
pollIntervalMs | number | 5000 | runAndWait のポーリング間隔 (ミリ秒) |
jobs | Record<string, JobConfig> | — | 公開する名前付き job。各キーが job アクセサーになります |
job単位の設定 (JobConfig)
| オプション | 型 | デフォルト | 説明 |
|---|---|---|---|
waitTimeout | number | 600000 | このjobのポーリングタイムアウトを上書きします |
taskType | TaskType | — | パラメータの自動マッピングに使用するタスクタイプ |
params | z.ZodType | — | 実行時のパラメータ検証に使用する Zod スキーマ |
環境変数
単一 job モード
DATABRICKS_JOB_ID を設定すると、default キー配下に job を 1 つ公開します:
DATABRICKS_JOB_ID=123456const handle = AppKit.jobs("default");マルチjobモード
job ごとに DATABRICKS_JOB_<NAME> を設定します:
DATABRICKS_JOB_ETL=123456
DATABRICKS_JOB_ML=789012const etl = AppKit.jobs("etl");
const ml = AppKit.jobs("ml");環境変数名は大文字、jobキーは小文字になります。環境から検出されたjobは、明示的な jobs 設定とマージされ、明示的な設定が優先されます。
パラメータの検証
params を使うと、実行時に Zod スキーマを強制できます。不正なパラメータは、job がトリガーされる前に 400 で拒否されます。
import { z } from "zod";
jobs({
jobs: {
etl: {
params: z.object({
startDate: z.string(),
endDate: z.string(),
dryRun: z.boolean().optional(),
}),
},
},
})タスクタイプのマッピング
taskType を設定すると、plugin は検証済みのパラメータを適切な SDK リクエストフィールドへ自動的にマッピングします。
| タスクタイプ | SDK フィールド | パラメータの形式 |
|---|---|---|
notebook | notebook_params | Record<string, string> — 値は文字列に変換されます |
python_wheel | python_named_params | Record<string, string> — 値は文字列に変換されます |
python_script | python_params | { args: string[] } — 位置引数 |
spark_jar | jar_params | { args: string[] } — 位置引数 |
sql | sql_params | Record<string, string> — 値は文字列に変換されます |
dbt | — | パラメータは受け付けません |
jobs({
jobs: {
etl: {
taskType: "notebook",
params: z.object({
startDate: z.string(),
endDate: z.string(),
}),
},
},
})taskType を省略した場合、パラメータはそのまま SDK に渡されます。
実行コンテキスト
jobは常にアプリのサービスプリンシパルとして実行されます。アプリのリソースバインディング (databricks.yml) でサービスプリンシパルに CAN_MANAGE_RUN が付与されるため、ユーザーは個別に権限を付与されなくてもrunを開始できます。Jobs UI のrunごとの実行者表示は、実際のユーザーではなくアプリのサービスプリンシパルになります。
HTTP endpoint
すべてのルートは /api/jobs 配下にマウントされます。
run をトリガーする
POST /api/jobs/:jobKey/run
Content-Type: application/json
{ "params": { "startDate": "2025-01-01" } }{ "runId": 12345 } を返します。
?stream=true を追加すると、run が完了するまでポーリングを行い、SSE でステータス更新を受信します:
POST /api/jobs/:jobKey/run?stream=true各 SSE イベントには { status, timestamp, run } が含まれます。
run の一覧取得
GET /api/jobs/:jobKey/runs?limit=20{ "runs": [...] } を返します。limit は 1〜100 の範囲に丸められ、デフォルトは 20 です。
run の詳細を取得する
GET /api/jobs/:jobKey/runs/:runId最新のステータスを取得
GET /api/jobs/:jobKey/status最新の run について { "status": "TERMINATED", "run": { ... } } を返します。
run をキャンセルする
DELETE /api/jobs/:jobKey/runs/:runId成功時は 204 No Content を返します。
プログラムによるアクセス
この plugin は、キーを指定して job を選択する呼び出し可能オブジェクトをエクスポートします:
const AppKit = await createApp({
plugins: [
server(),
jobs({
jobs: {
etl: { taskType: "notebook" },
},
}),
],
});
const etl = AppKit.jobs("etl");
// runをトリガーする
const result = await etl.runNow({ startDate: "2025-01-01" });
if (result.ok) {
console.log("Run ID:", result.data.run_id);
}
// runをトリガーして完了までポーリングする
for await (const status of etl.runAndWait({ startDate: "2025-01-01" })) {
console.log(status.status); // "PENDING"、"RUNNING"、"TERMINATED" など
}
// 読み取り操作
await etl.lastRun();
await etl.listRuns({ limit: 10 });
await etl.getRun(12345);
await etl.getRunOutput(12345);
await etl.getJob();
// キャンセル
await etl.cancelRun(12345);すべてのメソッドは ExecutionResult<T> を返します。result.data にアクセスする前に result.ok を確認してください。
実行のデフォルト値
| ティア | キャッシュ | リトライ | タイムアウト | メソッド |
|---|---|---|---|---|
| 読み取り | TTL 60秒 | 3回試行、1秒バックオフ | 30秒 | getRun、getJob、listRuns、lastRun、getRunOutput |
| 書き込み | 無効 | 無効 | 120秒 | runNow、cancelRun |
| ストリーム | 無効 | 無効 | 600秒 | runAndWait (SSE ポーリング) |