PipelineDocs
GitHub
Docs / はじめに / チュートリアル

チュートリアル — 手を動かして理解する

このページは、Pipeline を初めて触るエンジニア向けのハンズオンです。制御プレーンの起動から、GUI と REST API での最初のジョブ実行、自作プラグインの作成、そして「dispatcher → consumer」という実運用の基本形までを、順に手を動かして体験します。所要 20〜30 分です。

前提: Python 3.12+ と、セットアップで作った venv に 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 をそのまま返すだけ) を使います。

  1. 左メニュー ワークロード新規作成
  2. slug: echo-test / 表示名: Echoテスト
  3. 実行モード: python_module
  4. プラグイン: echo-sample / モジュール: main
  5. フォームに prefix / sleep_secs / fail_pk_substr が出ます (これが plugin.yaml の init_kwargs)。既定のまま 作成
  6. 一覧に echo-test が出たら、enabled トグルが ON になっていることを確認します。

次に 📨 投入 アイコンで、pk に hello、extra に {"foo": "bar"} を入れて投入します。数秒後、同じ行の 履歴 アイコンを開くと、run が緑 (success) で並び、展開すると output_jsonechoed が入っているのが見えます。

いま起きたこと: 投入で 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で解説しています。

なぜ2段に分けるのか: 「対象の発見 (I/O・低頻度)」と「重い処理 (GPU・高並列)」を別ワークロードにすると、後者だけを worker 追加で水平スケールでき、優先度も別々に付けられます。dispatcher は1台に固定 (max_concurrent_total: 1)、consumer は GPU ホストに複数、という配置が典型です。

次のステップ

Pipeline — GUI-first batch fleet · ぱっぷすラボ GitHub · ホーム