PipelineDocs
GitHub
Docs / プラグイン開発 / Storage SDK

Storage SDK

pipeline.storage は、各プラグインの setup() に散らばりがちな MinIO / DB 接続の生成・認証・死活監視を1箇所に集約するプラグイン向け SDK です。プラグインは接続を「名前」で参照し、認証情報をハードコードしません (Airflow Connections / Prefect Blocks に相当)。

公開 API

from pipeline.storage import StorageRegistry, MinioStore, Db, HealthStatus, load_env_file
シンボル役割
StorageRegistry名前付き接続を解決し、型付きラッパーを遅延生成 + キャッシュする接続レジストリ。
MinioStoreMinIO の型付きラッパー (リトライ・死活付き)。
DbMariaDB の型付きラッパー (再接続・死活付き)。
HealthStatus死活結果 (name / kind / ok / latency_ms / endpoint / error)。
load_env_file(path)KEY=VALUE.env を dict に読み込むヘルパ。

プラグインからの使い方

from pipeline.storage import StorageRegistry

def setup(**kwargs):
    # .env ファイル (creds) + init_kwargs overrides から接続レジストリを作る
    reg = StorageRegistry.from_env_file(kwargs.get("db_env_file"), overrides=kwargs)
    return {
        "reg": reg,
        "db":  reg.db(),           # MariaDB (name="main")
        "raw": reg.minio("raw"),   # MinIO 'raw' 接続
    }

def process(task, ctx, state):
    # MinIO: put/get/exists/remove (リトライ内蔵)
    state["raw"].put_file(key, tmp_path, content_type="image/jpeg")
    if not state["raw"].exists(key):
        return {"_error": "put failed"}
    # DB: query / execute
    rows = state["db"].query("SELECT ... WHERE id = %s", (task.pk,))
    return {"rows": len(rows)}

既知の接続名は minio("raw") / minio("crawl")db("main") です。エンドポイント・バケット・認証情報は .envMINIO_* / DB_* キーから解決されます (設定リファレンス)。未設定なら明確な例外を投げます。同じ (kind, name) は同一ラッパーがキャッシュ返却されます。

MinioStore のメソッド

メソッド説明
put_file(key, path, content_type=…)ファイルをアップロード。
get_to_file(key, path)オブジェクトをローカルへダウンロード。
exists(key)存在確認 (stat)。
remove(key)削除。

いずれも指数バックオフのリトライ (既定2回) で包まれます。

Db のメソッド

メソッド説明
query(sql, params)fetchall() の結果を返す。
execute(sql, params)INSERT/UPDATE 等を実行 (autocommit)。
ping()接続確認 (失敗時は接続を破棄し次回再接続)。

死活監視 (healthy / health)

# 単一接続
if not state["raw"].healthy():
    ...   # MinIO 落ち

# レジストリ一括 (監視ループ向け)
for h in reg.health():          # list[HealthStatus]
    if not h.ok:
        alert(f"{h.name} ({h.endpoint}) down: {h.error}")

MinioStore.healthy() は、minio パッケージがあれば認証付き bucket_exists で、無ければ /minio/health/live への素の HTTP GET で判定します。そのため minio ドライバの無い監視用 venv でも「接続拒否 / 到達不能 / タイムアウト」を検出できます。失敗時は壊れたクライアントを破棄し、次回再生成します。Db.healthy()ping()SELECT 1 です。

ドライバ (minio / mariadb) は各メソッド内で遅延 import されます。そのため制御プレーンやテスト用 venv にドライバが無くても、pipeline.storage の import 自体は失敗しません。
Pipeline — GUI-first batch fleet · ぱっぷすラボ GitHub · ホーム