チュートリアル — 手を動かして理解する
このページは、Pipeline を初めて触るエンジニア向けのハンズオンです。制御プレーンの起動から、GUI と REST API での最初のジョブ実行、自作プラグインの作成、そして「dispatcher → consumer」という実運用の基本形までを、順に手を動かして体験します。所要 20〜30 分です。
pip install -e . 済みであること。まだなら先にセットアップを済ませてください。ここでは単一プロセス mode (SQLite 1ファイル) を使うので、GPU も外部 DB も不要です。
ステップ0: 全体像を掴む
Pipeline の登場人物は3つだけです。まずこの関係を頭に入れると、以降の操作の意味が分かります。
- 制御プレーン — キューと Web UI/API を持つ司令塔。1プロセス。
- ワークロード — 「何をどう実行するか」の定義。専用のキュー (待ち行列) を1つ持ちます。
- ワーカー — キューからタスクを取り出して処理する実行役。制御プレーンに取りに行く pull 型です。
あなた / 外部システム / dispatcher
│ enqueue (タスク投入)
▼
┌───────────────┐ claim ┌──────────┐
│ ワークロード │ ───── タスク ─────▶ │ ワーカー │ → plugin.process()
│ のキュー │ ◀──── 結果 ──────── │ │
└───────────────┘ complete / fail └──────────┘
「タスクをキューに入れる」→「ワーカーが取り出して処理する」→「結果が run として残る」。この1往復がすべての基本です。
ステップ1: 起動して画面を見る
pipeline run --dev
次のようなバナーが出て、ワンタイムの管理者パスワードが表示されます。
Pipeline (dev) → http://localhost:8000
DB: sqlite:///./pipeline.db
admin password: xxxxxxxxxxxxxxxx (この起動限り)
ブラウザで http://localhost:8000 を開きます。ダッシュボードが表示され、「稼働ワーカー / 最近の失敗 / キュー深さ」の3パネルは空です (まだ何も無いので正常)。--dev は制御プレーンと同じプロセスの中にワーカーを1つ同居させるので、この後すぐタスクを処理できます。
ステップ2: GUI で最初のワークロードを作る
同梱の echo-sample プラグイン (受け取った pk と extra をそのまま返すだけ) を使います。
- 左メニュー ワークロード → 新規作成。
- slug:
echo-test/ 表示名:Echoテスト。 - 実行モード:
python_module。 - プラグイン:
echo-sample/ モジュール:main。 - フォームに
prefix/sleep_secs/fail_pk_substrが出ます (これが plugin.yaml のinit_kwargs)。既定のまま 作成。 - 一覧に
echo-testが出たら、enabled トグルが ON になっていることを確認します。
次に 📨 投入 アイコンで、pk に hello、extra に {"foo": "bar"} を入れて投入します。数秒後、同じ行の 履歴 アイコンを開くと、run が緑 (success) で並び、展開すると output_json に echoed が入っているのが見えます。
echo_test_queue テーブルに1行 (pending) が入り、同居ワーカーがそれを claim し、echo_sample.process() を呼び、戻り dict を runs に記録して行を消しました。これがライフサイクルの一巡です。
ステップ3: 同じことを REST API で
GUI でできることは全て HTTP API でもできます。CI やバッチ、外部システムからの投入はこちらを使います。
# タスクを1件投入
curl -X POST http://localhost:8000/api/v1/workloads/echo-test/tasks \
-H "Content-Type: application/json" \
-d '{"pk": "hello-2", "extra": {"foo": "bar"}}'
# 直近の run を確認
curl "http://localhost:8000/api/v1/workloads/echo-test/runs?limit=3"
# キューの深さ (pending / claimed / failed) を確認
curl http://localhost:8000/api/v1/workloads/echo-test/queue
より詳しい request/response は REST API リファレンスにまとめてあります。ブラウザで http://localhost:8000/docs を開けば Swagger UI で全エンドポイントを対話的に叩けます。
ステップ4: 自作プラグインを書く
ここが本番です。「URL を受け取り、HTTP GET してステータスとサイズを返す」プラグインを、空のディレクトリから作ります。外部依存を避けるため標準ライブラリ (urllib) だけで書きます。
4-1. ファイルを置く
プラグインルート (dev では ./plugins/、本番は PIPELINE_PLUGIN_ROOT = 既定 /opt/pipeline/plugins) に新しいディレクトリを作ります。
plugins/
url_check/
plugin.yaml
main.py
plugins/url_check/plugin.yaml
name: url-check
description: |
extra.url (無ければ pk) を HTTP GET し、ステータスコードとバイト数を返す。
タイムアウトやDNSエラーは例外にして Pipeline のリトライに任せる。
init_kwargs:
- key: timeout_secs
type: float
default: 10
label: HTTP timeout
min: 1
max: 120
- key: user_agent
type: str
default: pipeline-url-check/1.0
label: User-Agent
plugins/url_check/main.py
from __future__ import annotations
import time
import urllib.request
from typing import Any
def setup(**kwargs: Any) -> dict[str, Any]:
# プロセス起動時に1回。ここで作った dict が以降 state として渡る。
return {
"timeout": float(kwargs.get("timeout_secs", 10)),
"ua": str(kwargs.get("user_agent", "pipeline-url-check/1.0")),
}
def process(task: Any, ctx: Any, state: Any) -> dict[str, Any]:
# 1タスク = 1 URL。url は extra 優先、無ければ pk を URL とみなす。
url = (task.extra or {}).get("url") or task.pk
req = urllib.request.Request(url, headers={"User-Agent": state["ua"]})
t0 = time.monotonic()
# 例外 (タイムアウト/DNS/接続拒否) は握らずに送出 → fail 扱いになり自動リトライ。
with urllib.request.urlopen(req, timeout=state["timeout"]) as resp:
body = resp.read()
elapsed_ms = round((time.monotonic() - t0) * 1000)
# 4xx/5xx を「失敗」にしたいなら、例外を投げず _error を返すのが定石。
if resp.status >= 400:
return {"_error": f"HTTP {resp.status} for {url}"}
# 成功: dict を返すと runs.output_json に保存される。
return {"url": url, "status": resp.status, "bytes": len(body), "elapsed_ms": elapsed_ms}
def cleanup(state: Any) -> None:
# 今回は解放するものが無い。接続プールやファイルを持つならここで閉じる。
pass
4-2. ワークロードとして登録する
UI を再読み込みすると、プラグインドロップダウンに url-check が出ます (制御プレーンが plugins/ を走査するため)。ステップ2と同じ手順で、slug url-check / 実行モード python_module / プラグイン url-check / モジュール main で作成します。
4-3. 動かす
curl -X POST http://localhost:8000/api/v1/workloads/url-check/tasks \
-H "Content-Type: application/json" \
-d '{"pk": "https://example.com", "extra": {}}'
curl "http://localhost:8000/api/v1/workloads/url-check/runs?limit=1"
# → output_json: {"url":"https://example.com","status":200,"bytes":1256,"elapsed_ms":83}
わざと失敗させて、リトライと dead-letter も見てみましょう。到達できない URL を投げると、process() が例外を送出 → failed になり、max_attempts (既定5) まで自動で再試行され、上限で dead-letter として残ります。RunsDrawer で traceback が、Dashboard のキュー深さで failed 件数が確認できます。
process() の実体はワーカープロセス内にキャッシュされます (重い setup() を毎回走らせないため)。コードを直したらワーカーを再起動してください。dev mode なら pipeline run --dev を Ctrl-C して起動し直すだけ。本番は該当ホストの worker を restart します (運用Tips)。
ステップ5: プラグインをローカルで単体テストする
fleet に載せる前に、process() は普通の Python 関数として単体で叩けます。素早い反復に便利です。
# plugins/url_check/ で
python3 - <<'PY'
from types import SimpleNamespace
import main
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))
main.cleanup(state)
PY
task / ctx は本番では専用オブジェクトですが、process() が参照する属性 (pk / extra / attempt と、必要なら ctx.deadline 等) を持つダミーで十分です。属性一覧を参照してください。
ステップ6: 本物のパイプライン — dispatcher → consumer
ここまではタスクを手 (curl/UI) で入れていました。実運用では、タスクを生成し続ける仕掛けが要ります。Pipeline の定石は次の2段構成です。
- dispatcher — 上流 (フィード / DB / ディレクトリ / 外部 API) を定期的に見て、新しく処理すべき対象を下流ワークロードのキューへ
tasks/batchで流し込む。ループし続けるのでself_loop型で作ります。 - consumer — 流れてきたタスクを1件ずつ
process()する、ステップ4のような per-task プラグイン。
┌──────────┐ fetch ┌─────────────┐ tasks/batch ┌──────────────┐ process()
│ 上流ソース │ ─────────▶ │ dispatcher │ ─────────────▶ │ consumer │ ───────────▶ 結果
│ (feed/DB…) │ │ (self_loop) │ 下流キューへ │ (per-task) │
└──────────┘ └─────────────┘ enqueue └──────────────┘
dispatcher の心臓部は、下流キューへの batch enqueue と、自分の次 tick の self-enqueue です。
import json, urllib.request
def _post_json(url, payload):
req = urllib.request.Request(
url, data=json.dumps(payload).encode(),
headers={"Content-Type": "application/json"}, method="POST")
urllib.request.urlopen(req, timeout=15).read()
def process(task, ctx, state):
control = state["control_url"] # 例: http://localhost:8000
target = state["target_workload"] # 下流 consumer の slug
new_urls = discover_new_urls() # 上流を見て新規対象を集める (dedupは自前)
items = [{"pk": u, "extra": {}} for u in new_urls]
if items:
_post_json(f"{control}/api/v1/workloads/{target}/tasks/batch", {"items": items})
time.sleep(state["interval_s"]) # 次サイクルまで待つ
# 自分の次 tick を self-enqueue してループを継続 (これを忘れるとループが止まる)
_post_json(f"{control}/api/v1/workloads/{state['workload_slug']}/tasks",
{"pk": f"tick-{state['counter']+1}"})
return {"enqueued": len(items)}
実際に動く完全な実装は同梱の examples/plugins/rss_news_dispatcher/ (RSS を集めて rss-news-summarize に流す) が参考になります。self_loop の詳細・watchdog・singleton 化は プラグインSDKで解説しています。
max_concurrent_total: 1)、consumer は GPU ホストに複数、という配置が典型です。