コーディング × インフラ — asyncio GKE バッチワーカー改善

2026-04-25 AB: コーディング × インフラ ★★★☆☆ Python asyncio × GKE Autopilot SQLインジェクション対策 × Workload Identity

概要

🔒

SQLインジェクション対策

f-stringでのSQL埋め込みは危険。BigQueryのQueryJobConfig(query_parameters=...)でパラメータ化クエリを使う。

asyncio と同期I/Oの誤用

google-cloudクライアントは同期I/O。async defの中で直接呼んでもイベントループをブロック。run_in_executorでラップする。

🕐

タイムゾーン安全な日時

datetime.now()はタイムゾーンなし。本番ではdatetime.now(UTC)を使いisoformat()でRFC 3339形式に統一。

🧪

テスタビリティ向上(DI)

グローバルクライアントを廃止し、キーワード引数でクライアントを注入可能にする。テストでモッククライアントを渡せる。

問題

以下のコードの問題点を洗い出し、改善せよ。 このワーカーはGCP上のGKE AutopilotのPodで実行され、BigQueryからデータを取得して処理後にCloud Storageへ書き込む。

制約・前提条件

  • Python 3.12 + asyncio を使用
  • GKE Autopilot 上の Pod で実行
  • Workload Identity Federation(サービスアカウントキーファイルなし)
  • BigQuery・Cloud Storage のクライアントは google-cloud-python v3.x
  • batch_id は外部システムから渡ってくるユーザー入力値
期待する回答形式: 問題点の列挙 → 改善後コード(型ヒント・Google スタイル docstring・定数定義込み)→ 適用した設計パターン名

悪いコード (Before)

このコードには 6つの問題点 があります。見つけてみてください。
# worker.py(問題あり)
import asyncio
import json
from google.cloud import bigquery, storage
from datetime import datetime

BQ_CLIENT = bigquery.Client()  # グローバルに作成
GCS_CLIENT = storage.Client()

async def process_batch(batch_id):
    # BigQueryからデータ取得
    query = f"SELECT * FROM `my_project.my_dataset.orders` WHERE batch_id = '{batch_id}'"
    result = BQ_CLIENT.query(query).to_dataframe()

    # 処理
    processed = []
    for _, row in result.iterrows():
        processed.append({
            "order_id": row["order_id"],
            "total": row["amount"] * 1.1,  # 10%加算
            "processed_at": str(datetime.now())
        })

    # GCSに書き込み
    bucket = GCS_CLIENT.bucket("my-bucket")
    blob = bucket.blob(f"output/{batch_id}.json")
    blob.upload_from_string(json.dumps(processed))

    print(f"Done: {batch_id}")

async def main():
    batch_ids = ["batch_001", "batch_002", "batch_003"]
    tasks = [process_batch(bid) for bid in batch_ids]
    await asyncio.gather(*tasks)

if __name__ == "__main__":
    asyncio.run(main())

ヒント(段階的開示)

ヒント1 — 問題の4軸
問題は「セキュリティ」「非同期の誤用」「コードの保守性」「本番品質」の4軸から探す。
ヒント2 — 各問題のヒント
  • f"...WHERE batch_id = '{batch_id}'" → SQLインジェクションのリスク
  • BQ_CLIENT.query(query).to_dataframe() はブロッキング呼び出し → async def の中で使っても非同期にならない
  • datetime.now() はタイムゾーンなし → 本番では UTC 固定にすべき
  • グローバルクライアント → テスタビリティが低い、モジュールレベルの副作用
  • 1.1 というマジックナンバー → 定数化が必要
  • エラーハンドリングが皆無
ヒント3 — 修正方針

BigQuery の同期クライアントを asyncio で使うには:

loop = asyncio.get_running_loop()
result = await loop.run_in_executor(
    None, lambda: client.query(query).to_dataframe()
)

SQLインジェクション対策(BigQueryのパラメータクエリ):

job_config = bigquery.QueryJobConfig(
    query_parameters=[
        bigquery.ScalarQueryParameter("batch_id", "STRING", batch_id)
    ]
)
query = "SELECT * FROM `...` WHERE batch_id = @batch_id"

問題点分析

#問題点分類影響改善方法
1 f"WHERE batch_id = '{batch_id}'" でSQL直接埋め込み SQLインジェクション 外部入力からDBが操作される可能性 QueryJobConfig(query_parameters=[...]) でパラメータ化
2 async def 内でブロッキングI/O(BQ/GCS)を呼んでいる async誤用 イベントループをブロック。並行処理の効果がゼロ loop.run_in_executor(None, ...) でスレッドプールにオフロード
3 datetime.now() はタイムゾーンなし UTC問題 ローカル時刻でデータが汚染される datetime.now(UTC).isoformat() でRFC 3339形式に統一
4 グローバルクライアント(BQ_CLIENT, GCS_CLIENT テスタビリティ モジュールインポート時に副作用。テストで差し替え不可 ファクトリ関数 + キーワード引数でDI
5 1.1 というマジックナンバー 保守性 変更時に何箇所あるか分からない _SURCHARGE_RATE = 1.1 として定数化・コメント付与
6 エラーハンドリングが皆無(printのみ) 本番品質 失敗が無視され、データ欠損が発生しても気づかない try/except + logger.exception() + raise に変更

アーキテクチャ図 — GKE Autopilot + Workload Identity

GKE Autopilot Pod app container worker.py + run_in_executor KSA(K8s Service Account) Workload Identity バインド済み GSA(GCP Service Account) キーファイル不要 (WIF) BigQuery @batch_id パラメータ化 SQLインジェクション防止 GCP IAM roles/cloudsql.client のみ asyncio.gather 並列実行 Cloud Storage gs://my-bucket/output/{batch_id}.json run_in_executor でアップロード 3バッチを asyncio.gather で並列実行 → BQ同期I/OはThreadPoolExecutorでオフロード

模範解答

"""BigQuery orders バッチ処理ワーカー。

GKE Autopilot 上で動作し、Workload Identity Federation で
GCP リソースに認証する非同期バッチワーカー。
"""

from __future__ import annotations

import asyncio
import json
import logging
from datetime import UTC, datetime
from typing import Any

from google.cloud import bigquery, storage

# ─── 定数 ──────────────────────────────────────────────────────────────────
_PROJECT_ID = "my_project"
_DATASET = "my_dataset"
_TABLE = "orders"
_GCS_BUCKET = "my-bucket"
_GCS_OUTPUT_PREFIX = "output"
_SURCHARGE_RATE = 1.1  # 10% 加算レート(仕様書 §3.2 参照)

_QUERY = f"""
SELECT order_id, amount
FROM `{_PROJECT_ID}.{_DATASET}.{_TABLE}`
WHERE batch_id = @batch_id
"""

logger = logging.getLogger(__name__)


# ─── クライアントファクトリ(DI可能にする) ───────────────────────────────
def _make_bq_client() -> bigquery.Client:
    """BigQuery クライアントを生成する(Workload Identity 使用)。"""
    return bigquery.Client(project=_PROJECT_ID)


def _make_gcs_client() -> storage.Client:
    """Cloud Storage クライアントを生成する。"""
    return storage.Client(project=_PROJECT_ID)


# ─── 処理ロジック ────────────────────────────────────────────────────────
async def _fetch_orders(
    bq_client: bigquery.Client,
    batch_id: str,
) -> list[dict[str, Any]]:
    """BigQuery から指定バッチのオーダーを非同期取得する。

    Args:
        bq_client: BigQuery クライアント。
        batch_id: 処理対象バッチの識別子。

    Returns:
        オーダー行の辞書リスト。

    Raises:
        ValueError: batch_id が空の場合。
        google.api_core.exceptions.GoogleAPIError: BigQuery 呼び出し失敗時。
    """
    if not batch_id:
        raise ValueError("batch_id must not be empty")

    job_config = bigquery.QueryJobConfig(
        query_parameters=[
            bigquery.ScalarQueryParameter("batch_id", "STRING", batch_id)
        ]
    )

    loop = asyncio.get_running_loop()
    # BigQuery クライアントは同期 I/O → executor でブロッキングを回避
    df = await loop.run_in_executor(
        None,
        lambda: bq_client.query(_QUERY, job_config=job_config).to_dataframe(),
    )
    return df.to_dict(orient="records")


def _transform(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
    """オーダーに surcharge を付加して変換する。"""
    now_utc = datetime.now(UTC).isoformat()
    return [
        {
            "order_id": row["order_id"],
            "total": round(row["amount"] * _SURCHARGE_RATE, 2),
            "processed_at": now_utc,
        }
        for row in rows
    ]


async def _upload_result(
    gcs_client: storage.Client,
    batch_id: str,
    data: list[dict[str, Any]],
) -> None:
    """処理結果を GCS に非同期アップロードする。"""
    blob_name = f"{_GCS_OUTPUT_PREFIX}/{batch_id}.json"
    loop = asyncio.get_running_loop()
    await loop.run_in_executor(
        None,
        lambda: gcs_client.bucket(_GCS_BUCKET)
        .blob(blob_name)
        .upload_from_string(
            json.dumps(data, ensure_ascii=False),
            content_type="application/json",
        ),
    )
    logger.info("Uploaded gs://%s/%s (%d rows)", _GCS_BUCKET, blob_name, len(data))


async def process_batch(
    batch_id: str,
    *,
    bq_client: bigquery.Client | None = None,
    gcs_client: storage.Client | None = None,
) -> None:
    """1バッチを取得・変換・アップロードする(DI対応)。"""
    bq = bq_client or _make_bq_client()
    gcs = gcs_client or _make_gcs_client()

    logger.info("Processing batch: %s", batch_id)
    try:
        rows = await _fetch_orders(bq, batch_id)
        if not rows:
            logger.warning("No orders found for batch: %s", batch_id)
            return
        transformed = _transform(rows)
        await _upload_result(gcs, batch_id, transformed)
    except Exception:
        logger.exception("Failed to process batch: %s", batch_id)
        raise


async def main(batch_ids: list[str]) -> None:
    """複数バッチを並行処理する。"""
    await asyncio.gather(*(process_batch(bid) for bid in batch_ids))


if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    asyncio.run(main(["batch_001", "batch_002", "batch_003"]))

ポイント解説

1 SQLインジェクション対策(セキュリティ)
f-string で batch_id を直接埋め込むのは危険。BigQuery の QueryJobConfig(query_parameters=...) でパラメータ化クエリを使う。 @batch_id プレースホルダで値を安全にバインドする。
2 asyncio と同期 I/O の誤用修正
google-cloud-python の BigQuery・GCS クライアントは同期 I/O。 async def の中で直接呼んでも非同期にならず、イベントループをブロックする。 loop.run_in_executor(None, ...) でスレッドプールにオフロードして本来の並行性を得る。
3 タイムゾーン安全な日時
datetime.now() はローカル時刻でタイムゾーン情報なし。 本番では datetime.now(UTC) を使い、isoformat() で RFC 3339 形式に統一する。
4 テスタビリティ向上(DI パターン)
グローバルクライアントを廃止し、process_batch のキーワード引数でクライアントを注入可能にした。 テストでモッククライアントを渡せる。
✗ Before(テスト不可)
BQ_CLIENT = bigquery.Client()  # モジュール副作用
GCS_CLIENT = storage.Client()

async def process_batch(batch_id):
    # グローバルを直接使う → モック不可
    result = BQ_CLIENT.query(query)...
✓ After(DI対応)
async def process_batch(
    batch_id: str,
    *,
    bq_client=None,  # テストでモックを渡せる
    gcs_client=None,
) -> None:
    bq = bq_client or _make_bq_client()

実務への応用

GKE Autopilot + Workload Identity の構成

PodのService AccountにIAMバインドを付与するだけでクレデンシャルを自動取得できる。Terraform での設定例:

resource "google_service_account_iam_member" "workload_identity" {
  service_account_id = google_service_account.worker.name
  role               = "roles/iam.workloadIdentityUser"
  member             = "serviceAccount:${var.project_id}.svc.id.goog[${var.namespace}/${var.ksa_name}]"
}

Pod の serviceAccountName を KSA に設定するだけで google.cloud クライアントが自動的に WIF で認証する(キーファイル不要)。

次のステップ

発展問題: run_in_executor の代わりに asyncio-google-cloud ラッパーや vertexai.async 系を使った真のネイティブ非同期実装を検討する
  • 参考: Google Cloud Python クライアントライブラリの非同期サポートロードマップ
  • 参考: Python asyncio.Semaphore を使った並行数制限

今日のまとめ

BigQuery同期クライアントを asyncio 内で使うには run_in_executor でラップすること、 そして外部入力は必ずパラメータ化クエリで扱うこと—— この2点がセキュリティ・パフォーマンス両面で最重要。

グローバルクライアントをDI(依存性注入)に変えることで、 テストで差し替え可能になり、テスタビリティが大幅に向上する。

自己評価

自分の回答

気づき・メモ