プラグインSDK
プラグインは python_module executor から呼ばれる Python モジュールです。
3関数 (setup / process / cleanup) を export するのが基本契約で、
あとは plugin.yaml でフォーム項目とメタ情報を宣言するだけです。
ディレクトリ構成
plugins/
my_plugin/
plugin.yaml # マニフェスト (必須)
main.py # setup / process / cleanup を export (module名は任意)
requirements.txt # (任意) pip 依存
web/
panel.html # (任意) ui_panel: true のときのカスタムUI
README.md # (任意)
プラグインルートは環境変数 PIPELINE_PLUGIN_ROOT (既定 /opt/pipeline/plugins) です。ここに置いたディレクトリが GET /api/v1/plugins/available で列挙され、workload 作成フォームのプラグインドロップダウンに出ます。
plugin.yaml マニフェスト
name: my-plugin
description: |
ここにプラグインの説明を書きます。Web UI に表示されます。
# workload 作成フォームに出る入力項目 (setup() の kwargs になる)
init_kwargs:
- key: api_url
type: str
default: https://api.example.com
label: API endpoint
help: タスクが叩く API の URL
required: true
- key: timeout_secs
type: float
default: 10
label: HTTP timeout
min: 1
max: 300
- key: mode
type: enum
default: fast
label: 動作モード
options: [fast, safe]
- key: dry_run
type: bool
default: false
label: 書き込みをスキップ
- key: api_key
type: secret
label: API key
# フォームには出さず、raw JSON で編集する項目 (慣習: 制御プレーン URL 等)
hidden_kwargs:
- workload_slug
- control_url
# (任意) カスタム UI パネル
ui_panel: true
ui_panel_mode: image # video (既定) / image
# (任意) ループ型プラグインの宣言
workload_type: self_loop # 省略時は通常の per-task 型
マニフェストのトップレベルキー
| キー | 型 | 意味 |
|---|---|---|
name | str | プラグイン表示名。 |
description | str | Web UI に出す説明。 |
init_kwargs | list | フォーム項目。各値は setup() の kwargs になります。 |
hidden_kwargs | list[str] | フォームから隠す項目名 (raw JSON で直接編集する想定)。 |
ui_panel | bool | true で専用 UI パネルを有効化。 |
ui_panel_mode | str | パネルに渡す ?mode=。video (既定) / image。 |
workload_type | str/null | self_loop でループ型宣言 (watchdog 対象)。省略/event_driven は通常型。 |
init_kwargs エントリのフィールド
| フィールド | 型 | 意味 |
|---|---|---|
key | str (必須) | setup() に渡す kwarg 名。 |
type | str (必須) | int / float / str / path / bool / enum / secret。 |
default | any | フォームの初期値。 |
label / help | str | フォームのラベル / 補助説明。 |
min / max | number | 数値の下限 / 上限。 |
options | list | enum の選択肢。 |
required | bool | 必須フラグ。 |
type と Web UI の input 対応
| type | Web UI の input |
|---|---|
str / path | TextInput (path は path 用 placeholder) |
int / float | NumberInput (整数 / 小数) |
bool | Checkbox |
enum | Select (options 必須) |
secret | PasswordInput |
hidden_kwargs はあくまで「フォームに出さない」UI ヒントです。制御プレーンが workload_slug / control_url を自動で入れてくれる保証はないので、プラグイン側は必ずフォールバック付きで読みます (例: kwargs.get("control_url") or "http://localhost:8000")。実際に使う場合は workload の init_kwargs に raw JSON で値を入れてください。
main.py — プラグイン契約
from __future__ import annotations
from typing import Any
def setup(**kwargs: Any) -> Any:
"""プロセス起動時に1回だけ呼ばれます (最初の claim 時)。
重い初期化 (モデルロード / DB接続) はここで。戻り値が state になり、
process() / cleanup() に渡されます。省略可 (その場合 state=None)。
ここで例外を投げると executor build が失敗し、claim済タスクは失敗扱いに。"""
return {
"api_url": kwargs["api_url"],
"timeout": float(kwargs.get("timeout_secs", 10)),
"client": _make_http_client(),
}
def process(task: Any, ctx: Any, state: Any) -> dict[str, Any] | None:
"""1タスク毎に呼ばれます (必須)。
- task … pk / extra / attempt / workload_slug を持つ Task
- ctx … deadline / workdir / env / workload_config を持つ ExecutionContext
- state … setup() の戻り値 (キーワード引数 state= で渡される)
戻り dict は runs.output_json に保存されます。例外を上げると失敗扱いです。"""
r = state["client"].get(f"{state['api_url']}/{task.pk}", timeout=state["timeout"])
return {"status_code": r.status_code, "size": len(r.content)}
def cleanup(state: Any) -> None:
"""プロセス停止 / 設定変更 / filter変更 / eviction時に呼ばれます (任意)。
接続クローズなどを行います。ここでの例外は握り潰されます。"""
if state:
state["client"].close()
process(task, ctx, state=state) のように state をキーワード引数で渡します。第3引数の名前は必ず state にしてください。setup() は setup(**init_kwargs) で呼ばれるので、init_kwargs の各キーがそのままキーワード引数になります。
task / ctx / 戻り値
task (dataclass) — process() の第1引数。
| 属性 | 型 | 意味 |
|---|---|---|
pk | Any | タスクの主キー。 |
workload_slug | str | 所属 workload の slug。 |
attempt | int | 何回目の試行か (1始まり)。 |
extra | dict | enqueue 時に添えた JSON。 |
ctx (ExecutionContext) — 第2引数。
| 属性 | 型 | 意味 |
|---|---|---|
deadline | datetime | lease 失効時刻。長時間処理はこれを目安に切り上げます。 |
workdir | Path | タスク専用の一時ディレクトリ。処理後に自動削除されます。 |
env | dict | os.environ のコピー。 |
workload_config | dict | executor_config のコピー。 |
process() の戻り値の扱い
| 戻り値 | 結果 |
|---|---|
dict (キー _error なし) | 成功。output_json に保存。 |
dict で {"_error": "..."} を含む | 失敗 (error = その文字列)。例外を投げずに失敗を返せます。 |
None | 成功 (output_json なし)。 |
True / False | 成功フラグとしてそのまま使用。 |
例外を raise | 失敗 (traceback / stderr が記録され、attempt 増加)。 |
モジュール解決 (module / source_path / callable)
python_module の executor_config は3つの鍵でプラグインを特定します。
| 鍵 | 意味 |
|---|---|
source_path | プラグインのディレクトリ。dev は plugins/url_check、本番は /opt/pipeline/plugins/url_check 等。 |
module | そのディレクトリ内の .py ファイル名 (拡張子なし)。__init__.py を持つサブパッケージ名でも可。 |
callable | 呼ぶ関数名。既定 process。 |
内部では source_path が sys.path 先頭に加わり (from .lib import x のような相対 import が効く)、module は プラグイン slug 込みのユニークな sys.modules キーでロードされます。
dispatch_main.py のような同名ファイルを持つと、素の import では sys.modules キャッシュを共有し、片方のコードが両方の workload で動いてしまう事故が起きます (実際に発生しました)。Pipeline は source_path 指定時に spec_from_file_location + slug 入りキーでロードしてこれを物理的に隔離します。つまり ファイル名は他プラグインと被っても安全ですが、可読性のため main.py 等の分かりやすい名前を推奨します。
依存パッケージ (requirements.txt)
requirements.txt は「このプラグインが必要とする依存」を書くためのドキュメントです。Pipeline はこれを自動 install しません。プラグインが外部ライブラリを使うなら、ワーカーの venv に事前に入っている必要があります。方法は3つ。
- ワーカーの venv に手で
pip install -r requirements.txt。 - Deploy の
setup_commandにpip installを書いて配信時に流す。 - fleet 共通の venv を使い、依存をそこへ集約する。
標準ライブラリだけで書ければ依存ゼロで最も楽です (チュートリアルの url_check は urllib のみ)。
エラー処理とリトライ設計
- 失敗の表明は2通り。 一時的な障害 (ネットワーク・5xx・タイムアウト) は例外を
raiseしてmax_attemptsのリトライに乗せます。恒久的にダメな入力 (不正 URL・存在しない ID 等) は{"_error": "..."}を返して早々に確定させ、無駄なリトライを避けます。 - 冪等性を前提に書く。 lease 失効やリトライで、同じ pk が2回以上
process()されることがあります。出力の書き込みは upsert /INSERT IGNOREにする、成果物のキーに pk を使う等で、二重実行しても壊れないようにします。 - 時間を意識する。 長い処理は
ctx.deadlineを目安に切り上げます。max_duration_msを設定すると超過を自動で失敗にできます。 - OOM の自己申告。 エラー文言に CUDA OOM 系のシグネチャが含まれると、ワーカーが VRAM ピークを自己申告して以後の同時実行を絞ります。
ロギング
print(...) や logging の出力はワーカーの stdout/stderr に出て、ServiceLog 画面に集約されます。プラグインでは logging.getLogger("pipeline.plugin.url-check") のように名前を付けると追いやすくなります。例外は自動で traceback が runs と ServiceLog に残ります。
開発ループとローカルテスト
process() の実体はワーカープロセス内にキャッシュされます (重い setup() を毎回走らせないため)。したがってコードを直したらワーカーを再起動してください。dev は pipeline run --dev を Ctrl-C して起動し直すだけ、本番は該当ホストの worker を restart し、__pycache__/*.pyc を消してから上げ直します。
fleet に載せる前に、process() は普通の関数として単体で叩けます。
python3 - <<'PY'
from types import SimpleNamespace
import main # plugins/url_check/ で実行
state = main.setup(timeout_secs=5)
task = SimpleNamespace(pk="https://example.com", extra={}, attempt=1, workload_slug="url-check")
ctx = SimpleNamespace(deadline=None, workdir="/tmp", env={}, workload_config={})
print(main.process(task, ctx, state))
PY
手順を通しで体験するには チュートリアルのステップ4〜5 をどうぞ。
バッチ処理 (任意)
process_batch(tasks, ctx, state) を export すると、worker は batch_size に応じて複数タスクをまとめて渡します。戻り値の list 長は len(tasks) と一致させてください (不一致は全件失敗)。GPU 推論など、まとめた方が効率が良い plugin 向けです。
def process_batch(tasks: list, ctx: Any, state: Any) -> list[dict]:
embeddings = state["model"].encode([t.extra["text"] for t in tasks])
return [{"dim": len(e)} for e in embeddings] # len == len(tasks)
self_loop (ループ型) プラグイン
キューを消費する代わりに「回り続ける」プラグイン (dispatcher / 投入系 / supervisor) を作れます。workload_type: self_loop を宣言した上で、次の3点を実装します。
- ブートストラップ —
setup()の最後に、自分のキューへ最初の tick を HTTP 投入:
POST {control_url}/api/v1/workloads/{slug}/tasksに{"pk": "tick-1-<ts>"} - 長い tick —
process()の1回が実作業。ループ型はctx.deadline(= lease) の少し手前までwhileで回し、ポーリング型は1パス実行してtime.sleep(interval_s)。 - 次 tick の自己 enqueue —
process()の最後に、次の tick pk を再び自己投入してループを継続します。
worker はこれを通常の per-task plugin と同じく process(task, ctx, state) として駆動します (task は合成 tick)。
self_loop workload を走査し、pending=claimed=0 が300秒続いたらブートストラップ tick を1件投入)。また self_loop は多重起動すると増殖するため、max_concurrent_total: 1 + supervisor_enabled: false で単一化するのが定石です。
ストレージの使い方
MinIO / DB を使うプラグインは、接続の生成・認証・死活監視を Storage SDK (pipeline.storage) に任せられます。setup() で接続レジストリを作り、process() で名前指定のハンドルを使います。
サンプルと入手
最小サンプルが同梱されています。新規 plugin はこれをコピーするのが早いです。
examples/plugins/echo_sample/— pk/extra を echo するだけの最小 per-task plugin。examples/plugins/rss_news_dispatcher/— self_loop 型の分かりやすい実例。