AB: コーディング × インフラ複合 — BigQuery → GCS エクスポートバッチ改善

2026-04-18 (Day 6) AB: コーディング × インフラ ★★★☆☆ GKE / Argo Workflows / BigQuery 12-Factor App × ストリーミング処理

概要

🌐

12-Factor App (Config)

環境依存の設定は全て環境変数から取得。コードに設定値を埋め込まず、本番・ステージングを切り替え可能にする。

🌊

ストリーミング書き込み

blob.open("w")で行ごとにGCSへ書き込み、メモリ使用量をO(1)に抑える。数十万行でもOOMにならない。

💀

Fail Fast(sys.exit(1))

例外時に非ゼロで終了し、Argo Workflowsに失敗を伝える。終了コード0で終わるとArgoは成功と判断してしまう。

🛡️

パラメータ化クエリ

ScalarQueryParameterでSQLインジェクション・型ミスマッチを防ぐ。文字列結合でSQLを組み立てない。

問題

GKE上で動くPythonバッチジョブ。BigQueryから注文データを取得してCloud Storageにエクスポートするコードを改善してください。

前提条件

項目内容
実行環境GKE Autopilot上で Argo Workflows から呼ばれるバッチジョブ
データ量BigQueryのordersテーブルは日次で数十万行
認証方式Workload Identity Federationで認証(サービスアカウントキーファイルは不使用)
エラー処理エラー時はArgo Workflowsに失敗を伝える必要がある
環境本番・ステージング環境が存在する
期待する回答形式: (1) 問題点の列挙、(2) 改善後コード、(3) 実行例、(4) 適用した設計パターン名

悪いコード (Before)

このコードには 7つの問題 があります。「日付・エラー・環境・メモリ・ログ・SELECT *・SQLインジェクション」の観点から探してください。
# bad_export.py
import json
from google.cloud import bigquery, storage

def export():
    bq = bigquery.Client()
    gcs = storage.Client()

    result = bq.query("SELECT * FROM orders WHERE date = '2026-04-18'").result()

    data = []
    for row in result:
        data.append(dict(row))

    blob = gcs.bucket("my-bucket").blob("export.json")
    blob.upload_from_string(json.dumps(data))
    print("done")

export()

ヒント(段階的開示)

ヒント1 — 方向性
「日付ハードコード」「エラーハンドリング欠如」「環境依存のハードコード」「メモリ効率」「ログ」の5つの軸で問題がある。
ヒント2 — アプローチ
  • 日付: datetime.date.today() または環境変数・CLI引数で注入
  • エラー: try/except + sys.exit(1) でArgoに失敗を伝える
  • 環境: os.environ でバケット名・プロジェクトIDを注入(ハードコード禁止)
  • メモリ: 数十万行を一括 list() 化はOOM危険 → ストリーミング書き込み
  • ログ: print ではなく logging モジュールを使う
ヒント3 — コード骨格
import argparse, logging, os, sys
from datetime import date
from google.cloud import bigquery, storage

logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
logger = logging.getLogger(__name__)

BUCKET_NAME = os.environ["EXPORT_BUCKET"]
PROJECT_ID  = os.environ["GCP_PROJECT"]

def parse_args() -> argparse.Namespace: ...
def stream_export(target_date: date) -> None: ...

if __name__ == "__main__":
    args = parse_args()
    try:
        stream_export(args.date)
    except Exception as e:
        logger.error("Export failed: %s", e)
        sys.exit(1)

問題点分析

#問題点分類影響改善方法
1 日付のハードコード('2026-04-18'が固定) 保守性 毎日コードを書き換えないと使えない CLIフラグ --date or 環境変数で注入
2 エラーハンドリング欠如(例外でも終了コード0) 信頼性 Argo Workflowsは成功と判断してしまう sys.exit(1) で非ゼロ終了
3 環境依存のハードコード("my-bucket" 設定管理 本番・ステージングで同じコードを使えない 環境変数 EXPORT_BUCKET で切り替え
4 全行をメモリに蓄積(data = [] OOMリスク 数十万行でOut of Memory blob.open("w") でストリーミング書き込み
5 print でロギング 可観測性 Cloud Logging / DataDogで構造化ログが取れない logging モジュールを使う
6 SELECT * で全カラム取得 コスト 不要なカラムを全取得しコストとネットワーク量が増大 必要カラムのみ指定
7 文字列結合でSQL構築(SQLインジェクション的リスク) セキュリティ 型ミスマッチや不正な文字列でクエリが壊れる ScalarQueryParameter でパラメータ化

Before / After アーキテクチャ比較

✗ Before Argo Workflows 日付ハードコード export() 全行メモリに蓄積 BigQuery SELECT * GCS: my-bucket ⚠ ハードコード ✓ After Argo Workflows --date 注入 / env注入 stream_export() ストリーミング / logging BigQuery ScalarQueryParameter GCS: $EXPORT_BUCKET env変数で切り替え 例外発生時 sys.exit(1) → Argo失敗 問題: ハードコード / OOM / 終了コード0 改善: 12-Factor / ストリーミング / Fail Fast

模範解答

"""
orders_export.py — BigQuery → GCS エクスポートバッチ

設計:
- 環境変数で設定値を注入(12-Factor App)
- ストリーミング書き込みでメモリ効率化
- 構造化ロギング(DataDog / Cloud Logging 対応)
- sys.exit(1) で Argo Workflows に失敗を通知
"""

import argparse
import json
import logging
import os
import sys
from datetime import date, datetime

from google.cloud import bigquery, storage

# ---------------------------------------------------------------------------
# 定数(環境変数で上書き可能)
# ---------------------------------------------------------------------------
BUCKET_NAME: str = os.environ["EXPORT_BUCKET"]          # 例: "prod-export-bucket"
PROJECT_ID: str  = os.environ["GCP_PROJECT"]            # 例: "my-project-prod"
DATASET_ID: str  = os.environ.get("BQ_DATASET", "ecom") # デフォルト値あり

EXPORT_BLOB_TEMPLATE = "orders/date={date}/orders.jsonl"

# ---------------------------------------------------------------------------
# ロガー設定
# ---------------------------------------------------------------------------
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(name)s %(message)s",
)
logger = logging.getLogger(__name__)


def parse_args() -> argparse.Namespace:
    parser = argparse.ArgumentParser(description="BigQuery → GCS orders exporter")
    parser.add_argument(
        "--date",
        type=lambda s: datetime.strptime(s, "%Y-%m-%d").date(),
        default=date.today(),
        help="エクスポート対象日付 (YYYY-MM-DD). デフォルト: 本日",
    )
    return parser.parse_args()


def build_query(target_date: date) -> str:
    return f"""
        SELECT order_id, customer_id, amount, status, created_at
        FROM `{PROJECT_ID}.{DATASET_ID}.orders`
        WHERE DATE(created_at) = @target_date
    """


def stream_export(target_date: date) -> int:
    bq_client  = bigquery.Client(project=PROJECT_ID)
    gcs_client = storage.Client(project=PROJECT_ID)

    query = build_query(target_date)
    job_config = bigquery.QueryJobConfig(
        query_parameters=[
            bigquery.ScalarQueryParameter("target_date", "DATE", target_date.isoformat())
        ]
    )

    logger.info("Starting BigQuery query for date=%s", target_date)
    query_job = bq_client.query(query, job_config=job_config)

    blob_name = EXPORT_BLOB_TEMPLATE.format(date=target_date.isoformat())
    blob = gcs_client.bucket(BUCKET_NAME).blob(blob_name)

    row_count = 0
    # ストリーミング書き込み: 全行をメモリに乗せない
    with blob.open("w") as gcs_file:
        for row in query_job.result():
            gcs_file.write(json.dumps(dict(row), default=str) + "\n")
            row_count += 1
            if row_count % 10_000 == 0:
                logger.info("Streamed %d rows...", row_count)

    logger.info(
        "Export complete: date=%s rows=%d gs://%s/%s",
        target_date, row_count, BUCKET_NAME, blob_name,
    )
    return row_count


if __name__ == "__main__":
    args = parse_args()
    try:
        exported = stream_export(args.date)
    except Exception as exc:
        logger.error("Export failed: %s", exc, exc_info=True)
        sys.exit(1)  # Argo Workflows に非ゼロ終了コードで失敗を通知
# Argo Workflows step
- name: export-orders
  template: python-batch
  arguments:
    parameters:
    - name: export-date
      value: "{{workflow.parameters.date}}"

---
# Argo Template
env:
  - name: EXPORT_BUCKET
    valueFrom:
      secretKeyRef:
        name: export-config
        key: bucket
  - name: GCP_PROJECT
    value: "my-project-prod"
command: ["python", "orders_export.py", "--date", "{{inputs.parameters.export-date}}"]

# 実行例:
# export EXPORT_BUCKET="prod-orders-export"
# export GCP_PROJECT="my-project-prod"
# python orders_export.py --date 2026-04-18
# → 2026-04-18 09:00:01 INFO __main__ Starting BigQuery query for date=2026-04-18
# → 2026-04-18 09:00:05 INFO __main__ Streamed 10000 rows...
# → 2026-04-18 09:00:08 INFO __main__ Export complete: date=2026-04-18 rows=32481

ポイント解説

1 12-Factor App(Config)
環境依存の設定は全て環境変数から取得。コードに設定値を埋め込まない。本番・ステージング・開発環境を環境変数だけで切り替えられる。
2 ストリーミング処理
blob.open("w") で行ごとにGCSに書き込み、メモリ使用量を O(1) に抑える。数十万行でもメモリが枯渇しない。
3 パラメータ化クエリ
bigquery.ScalarQueryParameter でSQLインジェクション・型ミスマッチを防ぐ。文字列でSQLを組み立てない。
4 Fail Fast(sys.exit(1))
例外時に非ゼロで終了することでArgo Workflowsが失敗を検知できる。print("done")で終了コード0を返すと、失敗が上流に伝わらない。

設計パターン対応表

パターン適用箇所参考
12-Factor App (Config)環境変数で設定値を注入12factor.net
Fail Fastsys.exit(1) で即時失敗通知「良いコード・悪いコード」第1章
ストリーミング処理blob.open("w") 行ごとの書き込み「良いコード・悪いコード」第6章
パラメータ化クエリScalarQueryParameterOWASP Top 10 (A03: Injection)

実務への応用

Workload Identity Federation について

GKEノードのサービスアカウントに BigQuery/GCS 権限を付与する。コード内に認証情報は不要(Client() が自動で ADC を使用)。

  • サービスアカウントキーファイル(JSON)をコードに含めない
  • GKE Workload Identity で Pod のサービスアカウントに GCP リソースの権限を付与
  • Secretマネージャーと組み合わせて環境変数を安全に注入する

次のステップ

発展問題: エラー時にDead Letter Queue(Cloud Pub/Sub)に失敗イベントを送る処理を追加する
  • 参考キーワード: google.cloud.storage.Blob.open, bigquery.QueryJobConfig
  • 参考: Argo Workflows exit handler

今日のまとめ

「設定はコードに書かない・例外は必ず上流に伝える・大量データはストリーミング処理する」がGCP + K8sバッチジョブの3原則。

自己評価

自分の回答

気づき・メモ