概要
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(開発者の関心事)を分離することで、開発者は独立して自分のルートを変更できる。backendRefs の weight で10%カナリアを宣言的に実装できる。
問題 A: コーディング — asyncio.TaskGroup × ExceptionGroup × frozen dataclass(メール一括送信 Bad→Good)
以下の「悪いコード」は、ECサイト MOps チームのメール一括送信バッチです。問題点を全て洗い出し、asyncio.TaskGroup・ExceptionGroup / except*・frozen dataclass・Pydantic v2 バリデーション・型安全なエラーハンドリング を使って Bad→Good にリファクタリングしてください。
制約・前提条件
- Python 3.12+、asyncio、Pydantic v2、httpx を使うこと
asyncio.TaskGroupで並列送信し、except*で ExceptionGroup を捕捉すること- 送信結果を
frozen dataclassのBulkSendResultにまとめること - セマフォ(
asyncio.Semaphore)でレートリミット(同時10リクエスト)を制御すること - Google スタイル docstring・インラインコメント・名前付き定数を含めること
悪いコード (Before) — カテゴリ A
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: マジックナンバー(タイムアウト・並列数が即値)
asyncio.TaskGroup で並列化dict を raw で受け取る。Pydantic v2 EmailPayload で入力バリデーションexcept Exception as e: errors.append(str(e)) は全例外を文字列化するだけ。except* で型別に対応{"results": [], "errors": []} は型安全でない。frozen dataclass BulkSendResult で型安全にasyncio.Semaphore(10) で制御httpx.AsyncClient() を生成するとコネクションプールが使われない。async with httpx.AsyncClient() as client: で1つを共有MAX_CONCURRENT / SEND_TIMEOUT_SECONDS に切り出すヒント A(段階的開示)
ヒント1 — 方向性
asyncio.TaskGroup は Python 3.11+ で安定した並列タスク管理構造で、asyncio.gather(..., return_exceptions=True) に比べてキャンセル伝播が正確。except* は ExceptionGroup から特定の例外型だけを抜き出す Python 3.11+ の新構文。送信結果はミュータブルな dict より frozen dataclass の BulkSendResult で表現すると型安全でテストしやすい。
ヒント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 | 逐次送信 | テスト容易性 Ch11 | asyncio.TaskGroup で並列化 |
| 2 | 型ヒントなし・バリデーションなし | 型安全性 Ch2 | Pydantic v2 EmailPayload |
| 3 | 例外の握りつぶし | エラー処理 Ch10 | except* HTTPStatusError / TimeoutException |
| 4 | dict で結果を返す | 型安全性 Ch4 | frozen dataclass BulkSendResult |
| 5 | レートリミットなし | インフラ設計 | asyncio.Semaphore(MAX_CONCURRENT) |
| 6 | AsyncClient を毎回生成 | パフォーマンス | async with httpx.AsyncClient() as client: で共有 |
| 7 | マジックナンバー | 可読性 Ch7 | 名前付き定数 MAX_CONCURRENT / SEND_TIMEOUT_SECONDS |
模範解答 A
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}
"""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-api・segment-api・template-api)へのルーティングが 1つの Ingress に全部書かれていて、デプロイの粒度が荒い - TLS 証明書が手動管理(Secret に PEM を貼り付け)で、更新忘れリスクがある
HTTPRouteの重みルーティングによる カナリアデプロイが未実装
要件
| # | 要件 |
|---|---|
| 1 | GatewayClass(gke-l7-global-external-managed)と Gateway リソースを Terraform で定義すること |
| 2 | HTTPRoute で delivery-api(weight=90)・delivery-api-canary(weight=10)のカナリアルーティングを実装すること |
| 3 | TLS は Google マネージド証明書(ManagedCertificate / networking.gke.io/v1)を使い、自動更新にすること |
| 4 | segment-api と template-api は別の HTTPRoute リソースに分離し、独立してデプロイできる設計にすること |
ヒント B(段階的開示)
ヒント1 — 方向性
Gateway リソースがインフラ担当(クラスタ管理者)の関心事、HTTPRoute がアプリ担当(開発者)の関心事として責務を分離できる。GKE では GatewayClass: gke-l7-global-external-managed を指定すると Cloud Load Balancer が自動プロビジョニングされる。カナリアデプロイは HTTPRoute の backendRefs に weight フィールドで実装できる。
ヒント2 — Terraform のリソース構成
kubernetes_manifest.gatewayclass→GatewayClass(gke-l7-external)kubernetes_manifest.managed_cert→ManagedCertificate(networking.gke.io/v1)kubernetes_manifest.gateway→Gateway(HTTPS 443 + HTTP 80)kubernetes_manifest.route_delivery→HTTPRoute(delivery-api 90% + canary 10%)kubernetes_manifest.route_segment→HTTPRoute(/api/segment → segment-api)kubernetes_manifest.route_template→HTTPRoute(/api/template → template-api)Gateway listeners[].tls.options["networking.gke.io/pre-shared-certs"]で ManagedCertificate を参照
ヒント3 — HTTPRoute カナリアルーティングの骨格
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"
}
}
}]
}
}
}
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 ルーティング設計
模範解答 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
asyncio.gather(..., return_exceptions=True) は例外が起きても他タスクをキャンセルしない。TaskGroup は「いずれかが例外を送出したら残タスクを自動キャンセル」する。MOps の 2,400万通バッチで逐次送信を TaskGroup で並列化するだけでスループットが MAX_CONCURRENT 倍(10倍)以上になる。
except Exception は全例外を握りつぶす。except* httpx.HTTPStatusError は「HTTP ステータスエラーのみ」を捕捉し、except* httpx.TimeoutException は「タイムアウトのみ」を捕捉する。型別に処理を分けることで「HTTP 429 はリトライ、HTTP 400 はバリデーションエラーとして記録」という細かい対応が可能になる。
{"results": [], "errors": []} という dict は result["result"] というタイポで KeyError になる。@dataclass(frozen=True) の BulkSendResult は生成後に変更できないため意図しない副作用を防ぎ、result.success_rate と型安全にアクセスできる。@property で派生値を計算することで、テストが assert result.success_rate == 0.9 と読みやすく書ける。
カテゴリ B
旧 Ingress では全サービスのルーティングが1リソースに集中し、変更のたびにクラスタ管理者が介在する。
Gateway(インフラ担当: LB 設定・TLS 設定)と HTTPRoute(開発者担当: ルーティングルール)を分離することで、開発者は RBAC で許可された自分の HTTPRoute だけを独立して変更できる。Terraform でも route_segment だけ apply すれば gateway や route_delivery に影響しない。
旧 Ingress でのカナリアは
nginx.ingress.kubernetes.io/canary Annotation による Nginx 固有のハックが必要だった。Gateway API では backendRefs[].weight が K8s ネイティブな仕様で、GKE の Cloud LB が自動で重み付けルーティングを実装する。weight: 90/10 から 50/50 → 0/100 への段階的移行も terraform apply だけで完結する。
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-apiとdelivery-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 自動化と組み合わせることで「コードとしての宣言的インフラ」と「運用コストゼロの証明書管理」を同時に実現できる。