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 | 名前付き接続を解決し、型付きラッパーを遅延生成 + キャッシュする接続レジストリ。 |
MinioStore | MinIO の型付きラッパー (リトライ・死活付き)。 |
Db | MariaDB の型付きラッパー (再接続・死活付き)。 |
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") です。エンドポイント・バケット・認証情報は .env の MINIO_* / 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 自体は失敗しません。