A+B 複合 — asyncio.TaskGroup × except* ExceptionGroup × frozen dataclass(Ch4/Ch10/Ch11)× GKE Autopilot Gateway API + HTTPRoute カナリアデプロイ

2026-06-13 (Day 68) 土曜複合問題 ★★★★☆ Python 3.12 / asyncio / Pydantic v2 / frozen dataclass GKE Autopilot / Gateway API v1 / Terraform / ManagedCertificate

概要

asyncio.TaskGroup で並列送信・キャンセル伝播を正確に制御

asyncio.gather(..., return_exceptions=True) は例外が起きても他タスクをキャンセルしない。TaskGroup は「いずれかが例外を送出したら残タスクを自動キャンセル」し、except* で型別にエラーを処理できる。MOps の 2,400万通バッチで逐次送信を並列化するだけでスループットが10倍以上になる。

🔒

Semaphore でレートリミット — サーバー保護と SLA の両立

無制限の並列 HTTP は送信先 API を過負荷にする。asyncio.Semaphore(MAX_CONCURRENT) は「同時実行コルーチン数」を 10 に制限し、バックエンドの rate limit に合わせた安全な並列化を実現する。接続プール(httpx.AsyncClient を1つ共有)と組み合わせると効果が最大になる。

📦

frozen dataclass BulkSendResult — 値オブジェクトで型安全な結果

{"results": [], "errors": []} という dict は型安全でなく、result["result"] というタイポでランタイムエラーになる。@dataclass(frozen=True)BulkSendResult は生成後に変更できない値オブジェクトで、result.sent / result.success_rate と型安全にアクセスできる(Ch4/Ch2)。

🌐

Gateway API — Ingress の後継・責務分離でカナリアを宣言的に

Ingress は全サービスのルーティングが1リソースに集中し、変更のたびにクラスタ管理者が介在する。Gateway(管理者の関心事)と HTTPRoute(開発者の関心事)を分離することで、開発者は独立して自分のルートを変更できる。backendRefsweight で10%カナリアを宣言的に実装できる。

問題 A: コーディング — asyncio.TaskGroup × ExceptionGroup × frozen dataclass(メール一括送信 Bad→Good)

以下の「悪いコード」は、ECサイト MOps チームのメール一括送信バッチです。問題点を全て洗い出し、asyncio.TaskGroupExceptionGroup / except*frozen dataclassPydantic v2 バリデーション型安全なエラーハンドリング を使って Bad→Good にリファクタリングしてください。

制約・前提条件

  • Python 3.12+、asyncio、Pydantic v2、httpx を使うこと
  • asyncio.TaskGroup で並列送信し、except* で ExceptionGroup を捕捉すること
  • 送信結果を frozen dataclassBulkSendResult にまとめること
  • セマフォ(asyncio.Semaphore)でレートリミット(同時10リクエスト)を制御すること
  • Google スタイル docstring・インラインコメント・名前付き定数を含めること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード + 実行例(input→output)+ 適用した設計パターン名と書籍対応章

悪いコード (Before) — カテゴリ A

このコードには 7つの設計上の問題 が隠れています。見つけてみてください。
bad_bulk_sender.py — 逐次・型なし・例外握りつぶし
import asyncio
import httpx

async def send_all(emails, webhook_url):
    results = []
    errors = []

    # 問題1: 逐次送信(2,400万通に数時間かかる)
    for email in emails:
        try:
            # 問題2: 型ヒントなし・バリデーションなし
            resp = await httpx.AsyncClient().post(webhook_url, json=email)
            results.append(resp.status_code)
        except Exception as e:
            # 問題3: 例外を握りつぶす
            errors.append(str(e))

    # 問題4: dict で返す(型安全性なし)
    return {"results": results, "errors": errors}

# 問題5: レートリミットなし(サーバー過負荷)
# 問題6: httpx.AsyncClient を毎ループで生成(コネクション非効率)
# 問題7: マジックナンバー(タイムアウト・並列数が即値)
問題点サマリー(7点)
1逐次送信(Ch11) — for ループで1通ずつ送信。2,400万通では数時間かかる。asyncio.TaskGroup で並列化
2型ヒントなし・バリデーションなし(Ch2)dict を raw で受け取る。Pydantic v2 EmailPayload で入力バリデーション
3例外の握りつぶし(Ch10)except Exception as e: errors.append(str(e)) は全例外を文字列化するだけ。except* で型別に対応
4dict で結果を返す(Ch4){"results": [], "errors": []} は型安全でない。frozen dataclass BulkSendResult で型安全に
5レートリミットなし(インフラ設計) — 制限なし並列は送信先サーバーを過負荷にする。asyncio.Semaphore(10) で制御
6httpx.AsyncClient を毎回生成(パフォーマンス) — ループ内で httpx.AsyncClient() を生成するとコネクションプールが使われない。async with httpx.AsyncClient() as client: で1つを共有
7マジックナンバー(Ch7) — タイムアウト・並列数などが即値。名前付き定数 MAX_CONCURRENT / SEND_TIMEOUT_SECONDS に切り出す

ヒント A(段階的開示)

ヒント1 — 方向性
asyncio.TaskGroup は Python 3.11+ で安定した並列タスク管理構造で、asyncio.gather(..., return_exceptions=True) に比べてキャンセル伝播が正確。except* は ExceptionGroup から特定の例外型だけを抜き出す Python 3.11+ の新構文。送信結果はミュータブルな dict より frozen dataclassBulkSendResult で表現すると型安全でテストしやすい。
ヒント2 — アプローチ
  • asyncio.Semaphore(MAX_CONCURRENT) でレートリミット: async with semaphore: の中で HTTP リクエストを行う
  • httpx.AsyncClient はコンテキストマネージャで1つだけ生成し、全コルーチンで再利用(コネクションプール活用)
  • TaskGroup 内でタスクを生成 → except* httpx.HTTPStatusError で HTTP エラー、except* httpx.TimeoutException でタイムアウトを分けて処理
  • Pydantic v2 で EmailPayload を定義し、model_validate() で入力をバリデーション
  • @dataclass(frozen=True)BulkSendResult(sent=n, failed=m, errors=[...]) で結果を返す
ヒント3 — コードの骨格
MAX_CONCURRENT: int = 10
SEND_TIMEOUT_SECONDS: float = 30.0

@dataclass(frozen=True)
class BulkSendResult:
    sent: int
    failed: int
    errors: list[str]

    @property
    def success_rate(self) -> float:
        total = self.sent + self.failed
        return self.sent / total if total > 0 else 0.0

async def bulk_send(payloads: list[EmailPayload], webhook_url: str) -> BulkSendResult:
    semaphore = asyncio.Semaphore(MAX_CONCURRENT)
    sent: list[str] = []
    failed: list[str] = []

    async with httpx.AsyncClient(timeout=SEND_TIMEOUT_SECONDS) as client:
        try:
            async with asyncio.TaskGroup() as tg:
                for payload in payloads:
                    tg.create_task(_send_one(client, semaphore, payload, webhook_url, sent, failed))
        except* httpx.HTTPStatusError as eg:
            for exc in eg.exceptions:
                failed.append(str(exc))
        except* httpx.TimeoutException as eg:
            for exc in eg.exceptions:
                failed.append(f"Timeout: {exc}")

    return BulkSendResult(sent=len(sent), failed=len(failed), errors=failed)

問題点分析 — カテゴリ A

#問題点分類改善方法
1逐次送信テスト容易性 Ch11asyncio.TaskGroup で並列化
2型ヒントなし・バリデーションなし型安全性 Ch2Pydantic v2 EmailPayload
3例外の握りつぶしエラー処理 Ch10except* HTTPStatusError / TimeoutException
4dict で結果を返す型安全性 Ch4frozen dataclass BulkSendResult
5レートリミットなしインフラ設計asyncio.Semaphore(MAX_CONCURRENT)
6AsyncClient を毎回生成パフォーマンスasync with httpx.AsyncClient() as client: で共有
7マジックナンバー可読性 Ch7名前付き定数 MAX_CONCURRENT / SEND_TIMEOUT_SECONDS

模範解答 A

Before — 逐次・型なし・例外握りつぶし・dict 返し
import asyncio
import httpx

async def send_all(emails, webhook_url):
    results = []
    errors = []
    for email in emails:
        try:
            resp = await httpx.AsyncClient().post(webhook_url, json=email)
            results.append(resp.status_code)
        except Exception as e:
            errors.append(str(e))
    return {"results": results, "errors": errors}
After — TaskGroup・Semaphore・except*・frozen dataclass
"""bulk_sender.py — asyncio.TaskGroup × ExceptionGroup × frozen dataclass。

Ch2: 型の活用(Pydantic v2 EmailPayload)
Ch4: コレクション(frozen dataclass BulkSendResult)
Ch10: エラー処理(except* 型別ハンドリング)
Ch11: テスト容易性(TaskGroup + ファクトリ関数)
"""
from __future__ import annotations
import asyncio
import logging
from dataclasses import dataclass, field
from typing import Final

import httpx
from pydantic import BaseModel, EmailStr, field_validator

logger = logging.getLogger(__name__)

# 名前付き定数(マジックナンバー禁止)
MAX_CONCURRENT: Final[int] = 10            # 同時送信上限(レートリミット)
SEND_TIMEOUT_SECONDS: Final[float] = 30.0  # タイムアウト(秒)


class EmailPayload(BaseModel):
    """送信ペイロードの値オブジェクト(Pydantic v2 バリデーション付き)。"""
    recipient: EmailStr   # メールアドレス形式のバリデーション(Pydantic 組み込み)
    subject: str
    body_html: str
    campaign_id: str

    @field_validator("campaign_id")
    @classmethod
    def validate_campaign_id(cls, v: str) -> str:
        if not (4 <= len(v) <= 64):
            raise ValueError(f"campaign_id は 4〜64文字が必要。長さ: {len(v)}")
        return v


@dataclass(frozen=True)
class BulkSendResult:
    """一括送信の結果を保持するイミュータブル値オブジェクト。"""
    sent: int
    failed: int
    errors: list[str] = field(default_factory=list)  # frozen でも mutable field は OK

    @property
    def success_rate(self) -> float:
        """送信成功率(0.0〜1.0)。"""
        total = self.sent + self.failed
        return self.sent / total if total > 0 else 0.0

    @property
    def total(self) -> int:
        return self.sent + self.failed


async def _send_one(
    client: httpx.AsyncClient,
    semaphore: asyncio.Semaphore,
    payload: EmailPayload,
    webhook_url: str,
    sent_ids: list[str],   # 共有リスト(TaskGroup 内でミュータブルに追記)
    failed_ids: list[str],
) -> None:
    """単一メールを送信する内部コルーチン。"""
    async with semaphore:   # 同時実行数を MAX_CONCURRENT に制限
        resp = await client.post(webhook_url, json=payload.model_dump())
        resp.raise_for_status()  # 4xx/5xx → HTTPStatusError
    sent_ids.append(payload.campaign_id)
    logger.debug("送信成功: campaign_id=%s", payload.campaign_id)


async def bulk_send(
    raw_payloads: list[dict],
    webhook_url: str,
) -> BulkSendResult:
    """メールを並列で一括送信し、結果を BulkSendResult で返す。

    Returns:
        BulkSendResult: 送信成功数・失敗数・エラー一覧。
    """
    # Step 1: 入力バリデーション(不正データを早期に弾く)
    payloads: list[EmailPayload] = []
    validation_errors: list[str] = []
    for raw in raw_payloads:
        try:
            payloads.append(EmailPayload.model_validate(raw))
        except Exception as e:
            validation_errors.append(f"バリデーションエラー: {e}")

    semaphore = asyncio.Semaphore(MAX_CONCURRENT)
    sent_ids: list[str] = []
    http_errors: list[str] = []

    # Step 2: AsyncClient を1つだけ生成してコネクションプールを再利用
    async with httpx.AsyncClient(timeout=SEND_TIMEOUT_SECONDS) as client:
        try:
            # TaskGroup: 全タスク完了 or 例外発生で残タスクを自動キャンセル
            async with asyncio.TaskGroup() as tg:
                for payload in payloads:
                    tg.create_task(
                        _send_one(client, semaphore, payload, webhook_url, sent_ids, http_errors)
                    )
        except* httpx.HTTPStatusError as eg:
            # HTTP 4xx/5xx のみを捕捉(except* は型別フィルタ)
            for exc in eg.exceptions:
                http_errors.append(f"HTTP エラー {exc.response.status_code}: {exc}")
        except* httpx.TimeoutException as eg:
            # タイムアウトのみを捕捉
            for exc in eg.exceptions:
                http_errors.append(f"タイムアウト: {exc}")

    all_errors = validation_errors + http_errors
    return BulkSendResult(
        sent=len(sent_ids),
        failed=len(http_errors) + len(validation_errors),
        errors=all_errors,
    )
import asyncio

async def main():
    payloads = [
        {"recipient": "user1@example.com", "subject": "セール開始", "body_html": "<p>50%OFF</p>", "campaign_id": "camp-2026-summer"},
        {"recipient": "user2@example.com", "subject": "セール開始", "body_html": "<p>50%OFF</p>", "campaign_id": "camp-2026-summer"},
    ]
    result = await bulk_send(payloads, "https://hooks.example.com/send")

    print(result.sent)          # 2
    print(result.failed)        # 0
    print(result.success_rate)  # 1.0
    print(result.total)         # 2
    print(result.errors)        # []

    # バリデーションエラーのケース
    invalid = [{"recipient": "not-email", "subject": "X", "body_html": "", "campaign_id": "ab"}]
    r2 = await bulk_send(invalid, "https://hooks.example.com/send")
    print(r2.sent)    # 0
    print(r2.failed)  # 1  ← バリデーションエラーとしてカウント
    print(r2.errors)  # ["バリデーションエラー: ..."]

asyncio.run(main())
ポイント適用した設計原則/パターン書籍対応章
TaskGroup で並列実行・自動キャンセル伝播テスト容易性・非同期並列処理Ch11
Semaphore(10) でレートリミットインフラ保護・制御フロー設計Ch10
except* HTTPStatusError / TimeoutException型別エラーハンドリング・例外設計Ch10
@dataclass(frozen=True) BulkSendResult値オブジェクト・不変性Ch4 / Ch2
EmailPayload Pydantic v2 バリデーション型の活用・Fail-Fast 入力検証Ch2
AsyncClient を1つ共有コネクションプール再利用Ch11 パフォーマンス
MAX_CONCURRENT / SEND_TIMEOUT_SECONDSマジックナンバー排除・名前付き定数Ch7
# tests/test_bulk_sender.py
import asyncio
import pytest
from unittest.mock import AsyncMock, patch
from bulk_sender import bulk_send, BulkSendResult

class TestBulkSend:
    @pytest.mark.asyncio
    async def test_successful_send(self):
        """正常系: 全メール送信成功。"""
        payloads = [
            {"recipient": "u@example.com", "subject": "test", "body_html": "<p>hi</p>", "campaign_id": "camp-test"},
        ]
        mock_resp = AsyncMock()
        mock_resp.raise_for_status = AsyncMock()

        with patch("httpx.AsyncClient.post", return_value=mock_resp):
            result = await bulk_send(payloads, "https://example.com/hook")

        assert result.sent == 1
        assert result.failed == 0
        assert result.success_rate == 1.0
        assert isinstance(result, BulkSendResult)  # frozen dataclass

    @pytest.mark.asyncio
    async def test_validation_error(self):
        """バリデーションエラー: 不正メールアドレスで failed に計上。"""
        payloads = [
            {"recipient": "not-an-email", "subject": "x", "body_html": "", "campaign_id": "ab"},
        ]
        result = await bulk_send(payloads, "https://example.com/hook")

        assert result.sent == 0
        assert result.failed == 1  # バリデーションエラーが1件
        assert len(result.errors) == 1

    def test_bulk_send_result_is_frozen(self):
        """BulkSendResult が frozen(イミュータブル)であることを確認。"""
        r = BulkSendResult(sent=5, failed=1, errors=["err"])
        with pytest.raises(Exception):  # FrozenInstanceError
            r.sent = 999  # type: ignore[misc]

    def test_success_rate_zero_total(self):
        """合計ゼロの場合 success_rate は 0.0 を返す。"""
        r = BulkSendResult(sent=0, failed=0)
        assert r.success_rate == 0.0

問題 B: インフラ — GKE Autopilot × Kubernetes Gateway API × Terraform(MOps 配信 API ルーティング Bad→Good)

ECサイト MOps チームの配信 API は GKE Autopilot で稼働しています。以下の課題を解決してください。

現状の課題:
  • Ingress(networking.k8s.io/v1)で L7 ルーティングを設定しているが、Gateway API(gateway.networking.k8s.io/v1)への移行が K8s 1.30+ でデファクトになっている
  • 複数の配信サービス(delivery-apisegment-apitemplate-api)へのルーティングが 1つの Ingress に全部書かれていて、デプロイの粒度が荒い
  • TLS 証明書が手動管理(Secret に PEM を貼り付け)で、更新忘れリスクがある
  • HTTPRoute の重みルーティングによる カナリアデプロイが未実装

要件

#要件
1GatewayClassgke-l7-global-external-managed)と Gateway リソースを Terraform で定義すること
2HTTPRoutedelivery-api(weight=90)・delivery-api-canary(weight=10)のカナリアルーティングを実装すること
3TLS は Google マネージド証明書(ManagedCertificate / networking.gke.io/v1)を使い、自動更新にすること
4segment-apitemplate-api は別の HTTPRoute リソースに分離し、独立してデプロイできる設計にすること

ヒント B(段階的開示)

ヒント1 — 方向性
Kubernetes Gateway API は Ingress の後継。Gateway リソースがインフラ担当(クラスタ管理者)の関心事、HTTPRoute がアプリ担当(開発者)の関心事として責務を分離できる。GKE では GatewayClass: gke-l7-global-external-managed を指定すると Cloud Load Balancer が自動プロビジョニングされる。カナリアデプロイは HTTPRoutebackendRefsweight フィールドで実装できる。
ヒント2 — Terraform のリソース構成
  • kubernetes_manifest.gatewayclassGatewayClass(gke-l7-external)
  • kubernetes_manifest.managed_certManagedCertificate(networking.gke.io/v1)
  • kubernetes_manifest.gatewayGateway(HTTPS 443 + HTTP 80)
  • kubernetes_manifest.route_deliveryHTTPRoute(delivery-api 90% + canary 10%)
  • kubernetes_manifest.route_segmentHTTPRoute(/api/segment → segment-api)
  • kubernetes_manifest.route_templateHTTPRoute(/api/template → template-api)
  • Gateway listeners[].tls.options["networking.gke.io/pre-shared-certs"] で ManagedCertificate を参照
ヒント3 — HTTPRoute カナリアルーティングの骨格
Gateway 定義(Terraform kubernetes_manifest)
resource "kubernetes_manifest" "gateway" {
  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "Gateway"
    metadata = {
      name      = "mops-gateway"
      namespace = "mops"
    }
    spec = {
      gatewayClassName = "gke-l7-external"
      listeners = [{
        name     = "https"
        port     = 443
        protocol = "HTTPS"
        tls = {
          mode = "Terminate"
          options = {
            "networking.gke.io/pre-shared-certs" = "mops-managed-cert"
          }
        }
      }]
    }
  }
}
HTTPRoute カナリアルーティング(weight 90/10)
resource "kubernetes_manifest" "route_delivery" {
  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "HTTPRoute"
    ...
    spec = {
      rules = [{
        backendRefs = [
          {
            name   = "delivery-api"
            port   = 8080
            weight = 90   # 本番安定版
          },
          {
            name   = "delivery-api-canary"
            port   = 8080
            weight = 10   # カナリア版(新バージョン)
          }
        ]
      }]
    }
  }
}

アーキテクチャ図 — GKE Autopilot Gateway API ルーティング設計

Internet HTTPS requests Cloud Load Balancer(Gateway リソースが自動プロビジョニング) Gateway: mops-gateway GatewayClass: gke-l7-external ManagedCertificate api.mops.example.com(自動更新) GKE Autopilot クラスタ — mops namespace HTTPRoute: delivery-route match: PathPrefix /api/delivery delivery-api(安定版) weight: 90 → 90%のトラフィック delivery-api-canary weight: 10 → 10% weight: 90 + 10 = 100 → GKE が自動正規化 カナリア確認後: weight 100/0 → canary を stable に昇格 HTTPRoute: segment-route match: PathPrefix /api/segment → segment-api:8080 (weight: 100) HTTPRoute: template-route match: PathPrefix /api/template → template-api:8080 (weight: 100) Service: delivery-api port: 8080 Pod × N(GKE Autopilot 自動スケール) v1.2.0-stable Service: delivery-api-canary port: 8080 v1.3.0-canary(新バージョン) Service: segment-api port: 8080 / Pod × N Service: template-api port: 8080 / Pod × N Before: Ingress(networking.k8s.io/v1)— 廃止予定 全サービスのルーティングを1リソースに集中記述 / TLS は手動 Secret / カナリア不対応 / 開発者が管理者に依頼が必要 Terraform 管理領域(terraform/modules/gke-gateway/) kubernetes_manifest.gatewayclass kubernetes_manifest.managed_cert kubernetes_manifest.gateway kubernetes_manifest.route_* 責務分離: GatewayClass + Gateway(管理者)/ HTTPRoute(開発者)— 独立した terraform apply が可能 HTTPS /api/delivery /api/segment /api/template 90% 10%

模範解答 B

# terraform/modules/gke-gateway/main.tf
# GKE Autopilot Gateway API × ManagedCertificate × HTTPRoute カナリアデプロイ

locals {
  namespace = "mops"
  domain    = "api.mops.example.com"
}

# ── GatewayClass(GKE L7 外部ロードバランサを指定)──────────────────────────
resource "kubernetes_manifest" "gatewayclass" {
  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "GatewayClass"
    metadata = {
      name = "gke-l7-external"
    }
    spec = {
      # GKE コントローラー: このクラスを持つ Gateway に Cloud LB を自動プロビジョニング
      controllerName = "networking.gke.io/gateway"
    }
  }
}

# ── Google マネージド TLS 証明書(自動更新・手動管理不要)──────────────────
resource "kubernetes_manifest" "managed_cert" {
  manifest = {
    apiVersion = "networking.gke.io/v1"
    kind       = "ManagedCertificate"
    metadata = {
      name      = "mops-managed-cert"
      namespace = local.namespace
    }
    spec = {
      domains = [local.domain]  # Google CA が証明書を発行・90日ごとに自動更新
    }
  }
}

# ── Gateway(LB 設定: HTTPS 443 + HTTP 80 の2リスナー)────────────────────
resource "kubernetes_manifest" "gateway" {
  depends_on = [kubernetes_manifest.gatewayclass, kubernetes_manifest.managed_cert]

  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "Gateway"
    metadata = {
      name      = "mops-gateway"
      namespace = local.namespace
    }
    spec = {
      gatewayClassName = "gke-l7-external"
      listeners = [
        {
          name     = "https"
          port     = 443
          protocol = "HTTPS"
          tls = {
            mode = "Terminate"  # LB で TLS を終端(Pods へは HTTP で転送)
            options = {
              # GKE マネージド証明書を参照(PEM の手動貼り付け不要)
              "networking.gke.io/pre-shared-certs" = "mops-managed-cert"
            }
          }
        },
        {
          name     = "http"
          port     = 80
          protocol = "HTTP"
          # HTTP リクエストは HTTPRoute 側で HTTPS にリダイレクト
        }
      ]
    }
  }
}

# ── HTTPRoute: delivery-api カナリアルーティング(90/10 重み分散)──────────
resource "kubernetes_manifest" "route_delivery" {
  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "HTTPRoute"
    metadata = {
      name      = "delivery-route"
      namespace = local.namespace
    }
    spec = {
      parentRefs = [{
        name      = "mops-gateway"
        namespace = local.namespace
      }]
      hostnames = [local.domain]
      rules = [
        {
          matches = [{
            path = {
              type  = "PathPrefix"
              value = "/api/delivery"
            }
          }]
          backendRefs = [
            {
              name   = "delivery-api"        # 安定版(本番)
              port   = 8080
              weight = 90                    # 90%のトラフィックを安定版へ
            },
            {
              name   = "delivery-api-canary" # カナリア版(新バージョン検証中)
              port   = 8080
              weight = 10                    # 10%のトラフィックをカナリアへ
            }
          ]
        }
      ]
    }
  }
}

# ── HTTPRoute: segment-api(独立ルート・別チームが独立デプロイ可能)─────────
resource "kubernetes_manifest" "route_segment" {
  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "HTTPRoute"
    metadata = {
      name      = "segment-route"
      namespace = local.namespace
    }
    spec = {
      parentRefs = [{
        name      = "mops-gateway"
        namespace = local.namespace
      }]
      hostnames = [local.domain]
      rules = [
        {
          matches = [{
            path = {
              type  = "PathPrefix"
              value = "/api/segment"
            }
          }]
          backendRefs = [{
            name   = "segment-api"
            port   = 8080
            weight = 100  # 全トラフィックを segment-api へ(カナリアなし)
          }]
        }
      ]
    }
  }
}

# ── HTTPRoute: template-api(独立ルート)──────────────────────────────────
resource "kubernetes_manifest" "route_template" {
  manifest = {
    apiVersion = "gateway.networking.k8s.io/v1"
    kind       = "HTTPRoute"
    metadata = {
      name      = "template-route"
      namespace = local.namespace
    }
    spec = {
      parentRefs = [{
        name      = "mops-gateway"
        namespace = local.namespace
      }]
      hostnames = [local.domain]
      rules = [
        {
          matches = [{
            path = {
              type  = "PathPrefix"
              value = "/api/template"
            }
          }]
          backendRefs = [{
            name   = "template-api"
            port   = 8080
            weight = 100
          }]
        }
      ]
    }
  }
}

Bad vs Good 設計比較

観点Bad(旧 Ingress)Good(Gateway API)
ルーティング定義 1つの Ingress に全サービスを集中記述(1ファイル・1リソース) HTTPRoute を service ごとに分離(独立デプロイ・独立 terraform apply 可能)
TLS 管理 Secret に PEM を手動貼り付け(更新忘れ・有効期限管理が属人的) ManagedCertificate で Google CA が自動発行・自動更新(手動管理ゼロ)
カナリアデプロイ Ingress では重みルーティング非対応(Annotation ハックが必要) backendRefs.weight: 90/10 で宣言的カナリア(K8s ネイティブ)
責務分離 クラスタ管理者と開発者が同じ Ingress リソースを編集する Gateway(管理者)/ HTTPRoute(開発者)で RBAC と設計が一致
対応 K8s バージョン networking.k8s.io/v1(非推奨化進行中) gateway.networking.k8s.io/v1(K8s 1.30+ GA・デファクト)
ルーティングの粒度 変更するたびにクラスタ管理者の承認・作業が必要 HTTPRoute 単位で独立した apply → 影響範囲が限定される

カナリアデプロイの実施手順(weight 段階的移行)

【Step 1: カナリア起動(10%)】
  terraform apply → route_delivery: weight = [delivery-api: 90, canary: 10]
  kubectl get httproute delivery-route -n mops -o yaml  # 確認

【Step 2: DataDog で canary の観測(15〜30分)】
  - エラーレート: delivery-api vs delivery-api-canary を Service 別に監視
  - p99 レイテンシ: canary が stable の 1.2倍以内か確認
  - conversion rate: 10%カナリアのユーザーの CVR が落ちていないか確認

  判断基準:
    OK → Step 3 へ進む
    NG → ロールバック(weight: 100/0 に変更して terraform apply)

【Step 3: 段階的拡大(50%)】
  route_delivery: weight = [delivery-api: 50, delivery-api-canary: 50]
  terraform apply  # apply は HTTPRoute のみ変更(Gateway/GatewayClass は変更なし)

【Step 4: canary を stable に昇格(100%)】
  route_delivery: weight = [delivery-api: 0, delivery-api-canary: 100]
  terraform apply

  # または: canary を stable としてデプロイし直し、weight = [new-stable: 100]
  # delivery-api-canary の Deployment を delivery-api に昇格させる

【Step 5: canary Service を削除】
  kubectl delete service delivery-api-canary -n mops
  terraform apply  # route_delivery から canary backendRef を削除

【ロールバック(緊急時)】
  route_delivery: weight = [delivery-api: 100, delivery-api-canary: 0]
  terraform apply  # 数秒で完了(LB の設定変更のみ)
  # または: kubectl apply -f route_delivery_stable.yaml  # Terraform なしで即時対応

DataDog モニタリング設定(カナリア観測)

# DataDog モニター: カナリアのエラーレートが安定版の 2倍を超えたらアラート
monitors:
  canary_error_rate:
    type: metric alert
    query: >
      avg(last_5m):
        sum:trace.http.request.errors{service:delivery-api-canary,env:production}.as_rate()
        /
        sum:trace.http.request.hits{service:delivery-api-canary,env:production}.as_rate()
      > 0.02  # 2%以上のエラーレートでアラート
    message: |
      delivery-api-canary のエラーレートが閾値を超えました。
      ロールバックを検討してください:
        `kubectl apply -f k8s/route_delivery_stable.yaml`

ポイント解説

カテゴリ A

1 asyncio.TaskGroup でキャンセル伝播を正確に制御(Ch11)
asyncio.gather(..., return_exceptions=True) は例外が起きても他タスクをキャンセルしない。TaskGroup は「いずれかが例外を送出したら残タスクを自動キャンセル」する。MOps の 2,400万通バッチで逐次送信を TaskGroup で並列化するだけでスループットが MAX_CONCURRENT 倍(10倍)以上になる。
2 except* で型別エラーハンドリング(Ch10)
except Exception は全例外を握りつぶす。except* httpx.HTTPStatusError は「HTTP ステータスエラーのみ」を捕捉し、except* httpx.TimeoutException は「タイムアウトのみ」を捕捉する。型別に処理を分けることで「HTTP 429 はリトライ、HTTP 400 はバリデーションエラーとして記録」という細かい対応が可能になる。
3 frozen dataclass BulkSendResult で値オブジェクト(Ch4/Ch2)
{"results": [], "errors": []} という dict は result["result"] というタイポで KeyError になる。@dataclass(frozen=True)BulkSendResult は生成後に変更できないため意図しない副作用を防ぎ、result.success_rate と型安全にアクセスできる。@property で派生値を計算することで、テストが assert result.success_rate == 0.9 と読みやすく書ける。

カテゴリ B

4 Gateway API の責務分離 — 管理者と開発者の関心事を分ける
旧 Ingress では全サービスのルーティングが1リソースに集中し、変更のたびにクラスタ管理者が介在する。Gateway(インフラ担当: LB 設定・TLS 設定)と HTTPRoute(開発者担当: ルーティングルール)を分離することで、開発者は RBAC で許可された自分の HTTPRoute だけを独立して変更できる。Terraform でも route_segment だけ apply すれば gatewayroute_delivery に影響しない。
5 HTTPRoute weight で宣言的カナリアデプロイ
旧 Ingress でのカナリアは nginx.ingress.kubernetes.io/canary Annotation による Nginx 固有のハックが必要だった。Gateway API では backendRefs[].weight が K8s ネイティブな仕様で、GKE の Cloud LB が自動で重み付けルーティングを実装する。weight: 90/10 から 50/500/100 への段階的移行も terraform apply だけで完結する。
6 ManagedCertificate でTLS 運用コストをゼロに
Secret に PEM を手動貼り付けする設計は「有効期限管理・ローテーション・Key 漏洩リスク」が常につきまとう。GKE の ManagedCertificate は Google CA が証明書を自動発行・自動更新(Google 管理のサイクル)し、spec.domains にドメインを列挙するだけで複数ドメインの証明書も管理できる。cert-manager も不要で、運用の複雑さをゼロにできる。

実務への応用

  • TaskGroup × Semaphore は MOps バッチの基本パターン: Argo Workflows のステップから呼ばれる Python バッチで asyncio を使う場合、逐次 for ループから TaskGroup + Semaphore に変えるだけで送信スループットが大幅改善する。同時実行数は バックエンドの rate limit(例: メール送信 API が 10req/s)に合わせて MAX_CONCURRENT を調整する
  • frozen dataclass → pytest が読みやすくなる: BulkSendResult(sent=10, failed=2, errors=["..."]) のようにテストデータを直接コンストラクタで作れる。dict を {"sent": 10, "failed": 2} と書くより型安全で、result.success_rate == 0.83 と書けるアサーションは仕様として読める
  • Gateway API への移行は段階的に: 既存 Ingress と Gateway API は共存できる。まず新サービス(delivery-api-canary)を HTTPRoute で接続し、問題なければ既存サービスを順次移行する。kubectl get gateway -n mops でプロビジョニング状況を確認してから次のサービスを移行する
  • カナリアの観測は DataDog で Service 別に: delivery-apidelivery-api-canary のエラーレート・レイテンシ・スループットを DataDog のサービス別ダッシュボードで並べて監視。10%トラフィックでも統計的有意差が出たら即ロールバック(weight: 100/0 に変更して terraform apply
  • 証券マン視点: Gateway API 移行の ROI: Ingress 障害の平均対応時間(MTTR)が 2h → Gateway API 移行後は責務分離で問題が局所化し 30分に短縮。2名チームの時給 6,000円 × 年10回 × 1.5h = 18万円/年の削減。移行工数 3人日(18万円)は初年度で元が取れる

今日のまとめ

asyncio.TaskGroup + Semaphore + except* の組み合わせで「並列・レートリミット・型別エラーハンドリング」を宣言的に実現し、frozen dataclass BulkSendResult で結果を型安全な値オブジェクトとして返す設計が MOps バッチの基盤になる(Ch2/Ch4/Ch10/Ch11)。

インフラ側では GKE Gateway API の Gateway(管理者)/ HTTPRoute(開発者)による責務分離と weight カナリアルーティングが、K8s 1.30+ での分散システム設計のデファクトスタンダードであり、ManagedCertificate による TLS 自動化と組み合わせることで「コードとしての宣言的インフラ」と「運用コストゼロの証明書管理」を同時に実現できる。

自己評価

自分の回答

気づき・メモ