概要
asyncio.TaskGroup で逐次→並列(Python 3.11+)
for r in recipients: await send(r) の逐次処理は TaskGroup で全タスクを並列起動する。Semaphore=20 なら 1000件送信のレイテンシを約 1/20 に削減できる。失敗は ExceptionGroup でまとめて捕捉できる。
ContextVar で request_id をコンテキスト伝搬
全コルーチンの引数に request_id を追加する代わりに、ContextVar.set() で親コルーチンにセットするだけで全子コルーチンが ContextVar.get() で参照できる。asyncio のコンテキストコピー機構でスレッドセーフ。
structlog で JSON 構造化ログ(print 廃絶)
print(f"送信成功: {email}") は DataDog でのフィルタリング・アラート設定ができない。log.info("mail_sent", email=email, elapsed_ms=ms) の JSON ログにすることで DataDog Log Analytics での集計が可能になる。
frozen dataclass × except* で結果集約(Ch4/Ch10)
処理結果を BulkSendResult(frozen=True, slots=True) の不変値オブジェクトで返し、except* MailSendError(Python 3.11+ ExceptionGroup 構文)で型別に失敗を集約する。
問題
ECサイトの販促システムで、メール一括送信バッチ BulkMailSender がある。現在は逐次処理・ロギングが貧弱で、障害調査が困難になっている。以下の「悪いコード」には 7つの設計上の問題 がある。問題点を全て洗い出し、asyncio.TaskGroup / ContextVar / structlog / frozen dataclass を適用してリファクタリングせよ。
制約・前提条件
- Python 3.12+、
asyncio(TaskGroup、Semaphore)、httpx.AsyncClient、structlog(bindコンテキスト)を使うこと ContextVarを使い、request_idとcampaign_idをコルーチン間でスレッドセーフに共有すること- 同時送信数は
Semaphoreで上限管理すること(マジックナンバー禁止) - 送信失敗は例外を握りつぶさず、
MailSendErrorとして記録・集約すること - 処理結果は
@dataclass(frozen=True, slots=True)で返すこと - Google スタイル docstring・インラインコメント・名前付き定数を含めること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード(Google スタイル docstring・インラインコメント・名前付き定数含む)+ 実行例(input→output)+ 適用した設計パターン名と書籍の対応章
悪いコード (Before)
このコードには 7つの設計上の問題 が隠れています。見つけてみてください。
bad_bulk_mail_sender.py — 問題だらけのメール送信バッチ
import asyncio
import httpx
class BulkMailSender:
def __init__(self):
# 問題①: エンドポイントがマジックストリング(定数化なし)
self.endpoint = "https://mail-api.internal/v1/send"
# 問題②: self.endpoint が public(カプセル化なし)
async def send_bulk(self, campaign_id, recipients):
# 問題③: 逐次処理(1件ずつ await → 1000件で最大遅延)
results = []
for r in recipients:
result = await self._send_one(r, campaign_id)
results.append(result)
return results
async def _send_one(self, recipient, campaign_id):
email = recipient["email"]
async with httpx.AsyncClient() as client:
try:
resp = await client.post(
self.endpoint,
json={"to": email, "name": recipient.get("name", "")},
# 問題④: タイムアウト未設定(永遠に待ち続ける危険)
)
if resp.status_code == 200:
# 問題⑤: print でログを捨てる(structlog/logging 未使用)
print(f"送信成功: {email}")
return {"email": email, "status": "ok"}
else:
print(f"送信失敗: {email} - {resp.status_code}")
return {"email": email, "status": "failed"}
except Exception:
# 問題⑥: except Exception で例外を握りつぶす(詳細ログなし)
print(f"例外発生: {email}")
return {"email": email, "status": "error"}
# 問題⑦: request_id / campaign_id を引数で全コルーチンに渡す設計
# (ContextVar を使えばセット1回で全コルーチンが参照できる)
問題点サマリー(7点)
1エンドポイントがマジックストリング —
"https://mail-api.internal/v1/send" がハードコード。Final[str] 定数に抽出すること2
self.endpoint が public — 外部から上書き可能。self._endpoint でカプセル化(Ch7)3逐次処理(for ループの逐次 await) — 1000件 × 200ms = 200秒。
asyncio.TaskGroup + Semaphore で並列化すること4タイムアウト未設定 — 外部API障害時に永久待機。
httpx.AsyncClient(timeout=10.0) を設定すること5
print でログを捨てる — DataDog で集計不可。structlog の JSON 構造化ログに変換すること6
except Exception で握りつぶし — 障害調査不可。MailSendError を raise して呼び出し元で集約すること(Ch10)7
campaign_id を全コルーチンの引数に渡す — ContextVar で1箇所セットするだけで全コルーチンが参照できるヒント(段階的開示)
ヒント1 — 方向性
for ループの逐次 await は asyncio.TaskGroup で並列化できる。request_id を引数として全コルーチンに渡すのは煩雑 — ContextVar を使えば1箇所でセットして全コルーチンから参照できる。print / logging.info の構造化されていないログは structlog.get_logger().bind(...) で JSON ログに変換できる。タイムアウト未設定は httpx.AsyncClient(timeout=10.0) で解決する。
ヒント2 — アプローチ
asyncio.TaskGroup:async with asyncio.TaskGroup() as tg: tasks = [tg.create_task(self._send_one(client, sem, r)) for r in recipients] # 全タスク完了後に ExceptionGroup でまとめて例外が raise されるContextVar:REQUEST_ID_VAR: ContextVar[str] = ContextVar("request_id", default="") REQUEST_ID_VAR.set(run_id) # 親コルーチンで1回セット val = REQUEST_ID_VAR.get() # 子コルーチンから引数なしで参照structlog:log = structlog.get_logger().bind(request_id=REQUEST_ID_VAR.get(), email=email) log.info("mail_sent", elapsed_ms=elapsed) log.error("mail_failed", reason="timeout")except*(ExceptionGroup 捕捉):except* MailSendError as eg: for exc in eg.exceptions: errors.append(exc)
ヒント3 — コードの骨格
from contextvars import ContextVar
from typing import Final
import asyncio, structlog, httpx
from dataclasses import dataclass, field
MAX_CONCURRENT_SENDS: Final[int] = 20
HTTP_TIMEOUT_SECONDS: Final[float] = 10.0
MAIL_API_ENDPOINT: Final[str] = "https://mail-api.internal/v1/send"
REQUEST_ID_VAR: Final[ContextVar[str]] = ContextVar("request_id", default="")
CAMPAIGN_ID_VAR: Final[ContextVar[str]] = ContextVar("campaign_id", default="")
class MailSendError(RuntimeError):
def __init__(self, email: str, reason: str) -> None:
super().__init__(f"mail_send_failed: {email!r} — {reason}")
self.email = email
self.reason = reason
@dataclass(frozen=True, slots=True)
class BulkSendResult:
sent: int
failed: int
errors: tuple[MailSendError, ...] = field(default_factory=tuple)
class BulkMailSender:
def __init__(self, endpoint: str = MAIL_API_ENDPOINT) -> None:
self._endpoint = endpoint # カプセル化
async def send_bulk(self, campaign_id: str, recipients: list[dict]) -> BulkSendResult:
run_id = str(uuid.uuid4())
REQUEST_ID_VAR.set(run_id)
CAMPAIGN_ID_VAR.set(campaign_id)
sem = asyncio.Semaphore(MAX_CONCURRENT_SENDS)
errors: list[MailSendError] = []
async with httpx.AsyncClient(timeout=HTTP_TIMEOUT_SECONDS) as client:
try:
async with asyncio.TaskGroup() as tg:
for r in recipients:
tg.create_task(self._send_one(client, sem, r))
except* MailSendError as eg:
errors.extend(eg.exceptions)
return BulkSendResult(sent=len(recipients)-len(errors), failed=len(errors), errors=tuple(errors))
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | エンドポイントがマジックストリング | 可読性 | MAIL_API_ENDPOINT: Final[str] 定数に抽出 |
| 2 | self.endpoint が public — カプセル化なし | Ch7 コレクション隠蔽 | self._endpoint(アンダースコア)でカプセル化 |
| 3 | 逐次処理(for ループの逐次 await) | パフォーマンス | asyncio.TaskGroup + Semaphore で並列化 |
| 4 | タイムアウト未設定(永久待機の危険) | 堅牢性 | httpx.AsyncClient(timeout=HTTP_TIMEOUT_SECONDS) |
| 5 | print でログを捨てる(DataDog 集計不可) | 可観測性 | structlog で JSON 構造化ログに変換 |
| 6 | except Exception で握りつぶし(障害調査不可) | Ch10 エラー処理 | MailSendError を raise して呼び出し元で集約 |
| 7 | campaign_id を全コルーチン引数に渡す(煩雑) | 設計 | ContextVar で1箇所セット→全コルーチンが参照 |
模範解答
"""bulk_mail_sender.py — 非同期メール一括送信バッチ(Bad → Good リファクタリング)
良いコード・悪いコードで学ぶ設計入門(改訂新版)
- Ch4: 不変の活用(frozen dataclass BulkSendResult × tuple で不変コレクション)
- Ch7: コレクション操作の整理(_endpoint カプセル化 × errors list→tuple)
- Ch10: エラー処理(MailSendError 集約・握りつぶし禁止・except*)
asyncio.TaskGroup / ContextVar / structlog の実践パターン(2026年時点デファクト)
"""
from __future__ import annotations
import asyncio
import time
import uuid
from contextvars import ContextVar
from dataclasses import dataclass, field
from typing import Final
import httpx
import structlog
# ── 名前付き定数(マジックナンバー根絶)────────────────────────────────────
MAX_CONCURRENT_SENDS: Final[int] = 20 # 同時送信の上限(rate limit 対策)
HTTP_TIMEOUT_SECONDS: Final[float] = 10.0 # HTTPタイムアウト(永久待機防止)
MAIL_API_ENDPOINT: Final[str] = "https://mail-api.internal/v1/send"
# ── ContextVar でリクエストID・キャンペーンIDを非同期コンテキストに伝搬 ────
# asyncio はタスク生成時にコンテキストをコピーするため、
# 親コルーチンで set() すれば全子コルーチンが get() で参照できる
# → スレッドセーフ・引数不要・コルーチンの署名がシンプルになる
REQUEST_ID_VAR: Final[ContextVar[str]] = ContextVar("request_id", default="")
CAMPAIGN_ID_VAR: Final[ContextVar[str]] = ContextVar("campaign_id", default="")
# ── カスタム例外(Ch10: エラー処理)─────────────────────────────────────────
class MailSendError(RuntimeError):
"""単一宛先へのメール送信失敗を表す例外。
TaskGroup の ExceptionGroup からも型別捕捉できるよう
RuntimeError のサブクラスとして定義する。
Attributes:
email: 送信先メールアドレス。
reason: 失敗理由(HTTPステータスコードや "timeout" 等)。
"""
def __init__(self, email: str, reason: str) -> None:
super().__init__(f"mail_send_failed: {email!r} — {reason}")
self.email = email
self.reason = reason
# ── 送信結果値オブジェクト(Ch4: 不変の活用)──────────────────────────────
@dataclass(frozen=True, slots=True)
class BulkSendResult:
"""メール一括送信の処理結果を表す不変値オブジェクト。
Attributes:
sent: 送信成功件数。
failed: 送信失敗件数。
errors: 失敗した MailSendError のタプル(DLQ 用・不変コレクション)。
"""
sent: int
failed: int
# list ではなく tuple で不変コレクション(Ch7: コレクション操作の整理)
errors: tuple[MailSendError, ...] = field(default_factory=tuple)
# ── メール送信クラス ──────────────────────────────────────────────────────
class BulkMailSender:
"""非同期メール一括送信クラス。
asyncio.TaskGroup で並列送信し、ContextVar でリクエストIDを伝搬、
structlog で JSON 構造化ログを出力する。
Attributes:
_endpoint: メール送信 API エンドポイント(アンダースコアで隠蔽)。
"""
def __init__(self, endpoint: str = MAIL_API_ENDPOINT) -> None:
# アンダースコアでカプセル化(Ch7: コレクション操作の隠蔽)
self._endpoint = endpoint
async def send_bulk(
self,
campaign_id: str,
recipients: list[dict],
) -> BulkSendResult:
"""recipients 全員に asyncio.TaskGroup で並列メール送信する。
ContextVar にリクエストID・キャンペーンIDをセットすることで
全サブコルーチンが引数なしで参照できる。
Args:
campaign_id: 送信対象キャンペーンID(ContextVar に伝搬)。
recipients: 送信先リスト(各 dict に "email" / "name" キーを持つ)。
Returns:
送信成功件数・失敗件数・エラー詳細を含む BulkSendResult。
"""
if not recipients:
return BulkSendResult(sent=0, failed=0)
# ContextVar にセット(全サブコルーチンが get() で参照できる)
run_id = str(uuid.uuid4())
REQUEST_ID_VAR.set(run_id) # → 子コルーチンで REQUEST_ID_VAR.get()
CAMPAIGN_ID_VAR.set(campaign_id)
log = structlog.get_logger().bind(
request_id=run_id,
campaign_id=campaign_id,
total=len(recipients),
)
log.info("bulk_mail_started")
errors: list[MailSendError] = []
# Semaphore で同時送信数を MAX_CONCURRENT_SENDS に制限(rate limit 対策)
sem = asyncio.Semaphore(MAX_CONCURRENT_SENDS)
# httpx.AsyncClient を TaskGroup の外で生成して接続プールを再利用
async with httpx.AsyncClient(timeout=HTTP_TIMEOUT_SECONDS) as client:
try:
# TaskGroup: 全タスクを並列起動し、全完了を待つ
async with asyncio.TaskGroup() as tg:
for r in recipients:
tg.create_task(self._send_one(client, sem, r))
except* MailSendError as eg:
# except*(Python 3.11+): ExceptionGroup から型別に捕捉
# eg.exceptions は失敗した全 MailSendError のリスト
for exc in eg.exceptions:
errors.append(exc)
sent_count = len(recipients) - len(errors)
log.info("bulk_mail_finished", sent=sent_count, failed=len(errors))
return BulkSendResult(
sent=sent_count,
failed=len(errors),
errors=tuple(errors), # list → tuple で不変コレクション(Ch7)
)
async def _send_one(
self,
client: httpx.AsyncClient,
sem: asyncio.Semaphore,
recipient: dict,
) -> None:
"""単一宛先へのメール送信。
ContextVar からリクエストID・キャンペーンIDを取得して
structlog にバインドすることで、全ログに追跡IDが自動付与される。
Args:
client: 共有 httpx.AsyncClient(接続プールを再利用)。
sem: 同時送信数制限用 Semaphore(MAX_CONCURRENT_SENDS)。
recipient: 送信先 dict("email" / "name" キーを持つ)。
Raises:
MailSendError: 送信失敗(HTTP エラー・タイムアウト等)。
"""
email = recipient.get("email", "")
# ContextVar から参照(引数として渡す必要がない)
log = structlog.get_logger().bind(
request_id=REQUEST_ID_VAR.get(), # 親コルーチンが set() した値
campaign_id=CAMPAIGN_ID_VAR.get(),
email=email,
)
async with sem: # 同時実行数を MAX_CONCURRENT_SENDS に制限
start = time.monotonic()
try:
resp = await client.post(
self._endpoint,
json={"to": email, "name": recipient.get("name", "")},
)
resp.raise_for_status() # 4xx/5xx → httpx.HTTPStatusError
elapsed_ms = int((time.monotonic() - start) * 1000)
# structlog: JSON 構造化ログ(DataDog @elapsed_ms でフィルタ可)
log.info("mail_sent", elapsed_ms=elapsed_ms, status=resp.status_code)
except httpx.HTTPStatusError as exc:
log.error("mail_failed", reason=f"HTTP {exc.response.status_code}")
# 握りつぶし禁止: 呼び出し元(send_bulk)が except* で集約(Ch10)
raise MailSendError(email, f"HTTP {exc.response.status_code}") from exc
except httpx.TimeoutException:
log.error("mail_failed", reason="timeout")
raise MailSendError(email, "timeout")
正常系(3件並列送信)
import asyncio
async def main():
sender = BulkMailSender()
result = await sender.send_bulk(
campaign_id="CMP-2026-06",
recipients=[
{"email": "alice@example.com", "name": "Alice"},
{"email": "bob@example.com", "name": "Bob"},
{"email": "carol@example.com", "name": "Carol"},
],
)
print(f"sent={result.sent}, failed={result.failed}")
# sent=3, failed=0
asyncio.run(main())
structlog 出力(JSON 構造化ログ)
{"event": "bulk_mail_started", "request_id": "f3a1...", "campaign_id": "CMP-2026-06", "total": 3}
{"event": "mail_sent", "request_id": "f3a1...", "campaign_id": "CMP-2026-06", "email": "alice@example.com", "elapsed_ms": 142, "status": 200}
{"event": "mail_sent", "request_id": "f3a1...", "campaign_id": "CMP-2026-06", "email": "bob@example.com", "elapsed_ms": 158, "status": 200}
{"event": "mail_sent", "request_id": "f3a1...", "campaign_id": "CMP-2026-06", "email": "carol@example.com", "elapsed_ms": 137, "status": 200}
{"event": "bulk_mail_finished", "request_id": "f3a1...", "campaign_id": "CMP-2026-06", "sent": 3, "failed": 0}
異常系(2件失敗・1件成功)
result = await sender.send_bulk(
campaign_id="CMP-2026-06",
recipients=[
{"email": "alice@example.com"}, # 成功
{"email": "bad@example.com"}, # HTTP 429 (rate limited)
{"email": "slow@example.com"}, # timeout
],
)
print(f"sent={result.sent}, failed={result.failed}")
# sent=1, failed=2
for err in result.errors:
print(f" {err.email}: {err.reason}")
# bad@example.com: HTTP 429
# slow@example.com: timeout
# DataDog ログ(エラー行)
# {"event": "mail_failed", "email": "bad@example.com", "reason": "HTTP 429", ...}
# {"event": "mail_failed", "email": "slow@example.com", "reason": "timeout", ...}
パフォーマンス比較(1000件)
# Bad(逐次処理): 1000件 × 200ms = 200秒
# Good(TaskGroup + Semaphore=20): ceil(1000/20) × 200ms = 10秒
# → 約 20倍の高速化
# ContextVar の動作確認
import asyncio
from contextvars import ContextVar
VAR: ContextVar[str] = ContextVar("test", default="")
async def child():
print(VAR.get()) # → "set_by_parent"(引数なしで参照)
async def parent():
VAR.set("set_by_parent")
await child() # asyncio はタスク生成時にコンテキストをコピーする
asyncio.run(parent()) # → "set_by_parent"
| ポイント | 適用した設計原則/パターン | 書籍/仕様対応 |
|---|---|---|
asyncio.TaskGroup で並列化 | 並行タスク管理(PEP 654) | Python 3.11+ TaskGroup |
asyncio.Semaphore(MAX_CONCURRENT_SENDS) | rate limit 制御・マジックナンバー根絶 | Final 定数 + Semaphore |
ContextVar でリクエストID伝搬 | コンテキスト伝搬(引数不要化) | PEP 567 / asyncio context |
structlog.get_logger().bind(...) | JSON 構造化ログ(DataDog 対応) | structlog / OTel 1.x |
except* MailSendError as eg | ExceptionGroup 型別捕捉(Ch10) | PEP 654 / Python 3.11+ |
BulkSendResult(frozen=True, slots=True) | 不変値オブジェクト(Value Object) | Ch4 不変の活用 |
errors: tuple[MailSendError, ...] | 不変コレクション(list 禁止) | Ch7 コレクション操作の整理 |
self._endpoint(アンダースコア) | クラスによるカプセル化 | Ch7 コレクション操作の隠蔽 |
設計図 — Bad(逐次)vs Good(並列 + ContextVar)
ポイント解説
1
asyncio.TaskGroup で逐次→並列(Python 3.11+)for r in recipients: await send(r) の逐次処理は、1000件×200ms = 200秒かかる。TaskGroup で全タスクを並列起動し、Semaphore(20) で同時実行数を制御することで、ceil(1000/20) × 200ms = 10秒に短縮できる。TaskGroup はすべてのタスク完了を待ち、失敗は ExceptionGroup でまとめて伝搬される(asyncio.gather よりも安全なエラー伝搬)。
2
ContextVar でコンテキスト伝搬(引数不要化)request_id を全コルーチンの引数に追加すると、関数シグネチャが複雑になる。ContextVar を使えば、親コルーチンで REQUEST_ID_VAR.set(run_id) の1回だけで、全子コルーチンが REQUEST_ID_VAR.get() で参照できる。asyncio はタスク生成時にコンテキストをコピーするため、スレッドセーフかつ他のタスクに影響しない(contextvars.copy_context() が内部で使われる)。
3
structlog で JSON 構造化ログ(DataDog 対応)print(f"送信成功: {email}") は DataDog での @email フィールドによる絞り込みや、キャンペーン別成功率の集計ができない。structlog.get_logger().bind(request_id=..., email=...).info("mail_sent", elapsed_ms=...) の形式にすることで、DataDog Log Analytics で @campaign_id 別の送信成功率・レイテンシを1クエリで集計できる。
4
except* で ExceptionGroup を型別捕捉(Python 3.11+)asyncio.TaskGroup で複数タスクが失敗すると ExceptionGroup が raise される。except* MailSendError as eg で型別に捕捉し、eg.exceptions から個別エラーを取り出して errors リストに集約する。except* は ExceptionGroup の中から指定した型のみを取り出す構文で、他の型の例外は上位に再伝搬される。
5
外部メール API の rate limit(例: 20 requests/sec)に引っかからないよう、
asyncio.Semaphore で同時送信数を制限(rate limit 対策)外部メール API の rate limit(例: 20 requests/sec)に引っかからないよう、
asyncio.Semaphore(MAX_CONCURRENT_SENDS) の async with sem: ブロックで同時実行数を制御する。MAX_CONCURRENT_SENDS: Final[int] = 20 の定数化により、マジックナンバーを根絶し、環境に応じて1箇所で変更できる。
6
処理結果を不変値オブジェクトで返すことで、呼び出し元が
BulkSendResult(frozen=True, slots=True) 値オブジェクト(Ch4/Ch7)処理結果を不変値オブジェクトで返すことで、呼び出し元が
result.sent = 999 のような誤変更を行えなくなる(Ch4)。errors: tuple[MailSendError, ...] で不変コレクションにすることで、呼び出し元が .append() でエラーを追加するような副作用を防ぐ(Ch7)。
実務への応用
- Argo Workflows メール送信ステップ:
BulkSendResult.failed > 0の場合に DataDog カスタムメトリクスを送信し、Slack アラートを発火するパターンに直結する。result.errorsを Dead Letter Queue(Pub/Sub トピック)に送ることで再処理も可能。 - MOps クーポン一括送信の可観測性:
CAMPAIGN_ID_VARをstructlogにバインドすることで、DataDog Log Analytics で@campaign_id別の送信成功率・平均レイテンシを1クエリで集計できる。DataDog APM との相関ログにはstructlogプロセッサでtrace_id/span_idを追加する。 - OpenTelemetry との統合:
structlogの出力プロセッサに OTLP エクスポータを追加することで、DataDog のネイティブ OTLP ingestion(Agent v7.61+)との連携が可能になる。request_idを OTel のtrace_idとして使うことで、Argo ジョブログと APM トレースを相関させられる。 - GKE Autopilot + Spot Node 対策: Spot Node のプリエンプション時に送信中のバッチが途切れる。
BulkSendResult.failed_rowsを Cloud Spanner / Pub/Sub に保存しておくことで、再起動後のリトライが可能になる。
今日のまとめ
asyncio.TaskGroup × Semaphore で逐次→並列(約20倍高速化)、ContextVar で request_id / campaign_id をコルーチン間で引数なしに伝搬、structlog の JSON 構造化ログで DataDog での可観測性を確立する。except*(ExceptionGroup)で型別に失敗を集約し、BulkSendResult(frozen=True) で不変値オブジェクトとして返す設計は、Argo Workflows / MOps バッチの標準パターンとして押さえておくこと。