概要
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
模範解答
"""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 で
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 内で使うには
グローバルクライアントをDI(依存性注入)に変えることで、 テストで差し替え可能になり、テスタビリティが大幅に向上する。
run_in_executor でラップすること、
そして外部入力は必ずパラメータ化クエリで扱うこと——
この2点がセキュリティ・パフォーマンス両面で最重要。グローバルクライアントをDI(依存性注入)に変えることで、 テストで差し替え可能になり、テスタビリティが大幅に向上する。