概要
asyncio.Queue Bounded Buffer でバックプレッシャー(Ch11)
asyncio.Queue(maxsize=MAX_QUEUE_SIZE) は put() が満杯時に自動 await し、プロデューサーを下流ワーカーの処理速度に合わせて減速させる。無制限キューによるメモリ無限増大を防ぐ唯一の簡潔な仕組み。
asyncio.TaskGroup で構造化並行性(Ch10/Ch11)
async with asyncio.TaskGroup() as tg: 内のタスクが1つでも例外を投げると、残りの全タスクを自動キャンセルして ExceptionGroup として再送出する。gather(return_exceptions=True) と異なり、エラーが黙殺されない。
contextvars.ContextVar で暗黙的コンテキスト伝播(Ch10)
campaign_id を全コルーチンの引数に渡すと引数が肥大化する。ContextVar.set() で設定し ContextVar.get() で取得することで、ロギング・トレース文脈を自然に伝播。token を finally でリセットし汚染を防ぐ。
PEP 695 type 文 × frozen dataclass(Ch2/Ch4)
type JobItem = tuple[str, str] の PEP 695 型エイリアスは TypeAlias より宣言的。@dataclass(frozen=True, slots=True) DispatchResult は値オブジェクトとして書き換え禁止・軽量化。failed_ids: tuple[str, ...] で hashable かつ immutable を保証。
問題
ECサイト MOps チームでは、Argo Workflows から起動されるキャンペーン配信バッチを Python で実装している。以下の「悪いコード」はワーカープールを threading.Thread + list で自前管理し、バックプレッシャー制御なしでメッセージを無制限にエンキュー、contextvars.ContextVar を未使用、型エイリアスもマジックストリングで散在している。
問題点を全て洗い出し、asyncio.Queue(Bounded Buffer)・asyncio.TaskGroup・contextvars.ContextVar・type 文(PEP 695)・@dataclass(frozen=True, slots=True) を使って Bad→Good にリファクタリングしてください。
制約・前提条件
- Python 3.12+(PEP 695
type文使用可) asyncio.Queue(maxsize=N)でバックプレッシャー制御(プロデューサーがawait queue.put()でブロック)asyncio.TaskGroupでワーカーをまとめ、例外伝播を構造化並行性で保証することcontextvars.ContextVarでcampaign_idをコルーチン間の暗黙的コンテキストとして伝播することtype JobItem = tuple[str, str]の PEP 695 型エイリアスを使うこと- 送信結果は
@dataclass(frozen=True, slots=True)のDispatchResultで返すこと - Google スタイル docstring・インラインコメント・名前付き定数を含めること
悪いコード (Before)
import queue
import threading
import logging
import time
logger = logging.getLogger(__name__)
# 問題1: マジックナンバー散在 — 100, 5 がコードに埋め込まれている
# 問題2: 型エイリアスなし — (str, str) が何を意味するか不明
JOBS: list = [] # 型情報ゼロ
def send_message(campaign_id: str, recipient_id: str, message: str) -> dict:
# 問題3: 返り値が dict — 型安全性ゼロ、キーのタイポに気づけない
time.sleep(0.01) # 外部 API 呼び出しのシミュレート
import random
if random.random() < 0.05:
return {"status": "failed", "id": recipient_id}
return {"status": "success", "id": recipient_id}
def worker_thread(job_queue, campaign_id, results, lock):
# 問題4: campaign_id を全関数に引数として渡す(引数バケツリレー)
failed = []
sent = 0
while True:
try:
item = job_queue.get(timeout=1)
except Exception:
break
recipient_id, message = item
sent += 1
result = send_message(campaign_id, recipient_id, message)
if result["status"] == "failed": # 問題5: マジックストリング "failed"
failed.append(recipient_id)
job_queue.task_done()
with lock:
results.append({"sent": sent, "failed": failed}) # 問題6: dict + mutable list
def run_campaign(campaign_id: str, jobs: list) -> dict:
# 問題7: queue.Queue() — maxsize 未指定で無制限(バックプレッシャーなし)
job_queue = queue.Queue() # メモリ無限増大の原因
results = []
lock = threading.Lock()
threads = []
# 全ジョブを一気に投入 — プロデューサー/コンシューマー分離なし
for job in jobs:
job_queue.put(job)
# 問題8: threading.Thread — GIL 制約・キャンセル不可・例外伝播が困難
for _ in range(5): # マジックナンバー 5
t = threading.Thread(
target=worker_thread,
args=(job_queue, campaign_id, results, lock)
)
t.start()
threads.append(t)
job_queue.join()
for t in threads:
t.join()
total_sent = sum(r["sent"] for r in results)
total_failed = sum(len(r["failed"]) for r in results)
return {
"total_sent": total_sent,
"total_success": total_sent - total_failed,
"total_failed": total_failed,
}
100, 5, 0.05 がコードに埋め込まれている。MAX_QUEUE_SIZE, NUM_WORKERS, SUCCESS_RATE_THRESHOLD で名前付き定数にlist が何のリストか不明。type JobItem = tuple[str, str] の PEP 695 型エイリアスで意図を明示frozen dataclass DispatchResult で型安全な値オブジェクトにcontextvars.ContextVar で暗黙的に伝播StrEnum JobStatus で型安全に管理frozen=True, slots=True DispatchResult + tuple[str, ...] でイミュータブル化asyncio.Queue(maxsize=MAX_QUEUE_SIZE) で制御asyncio.TaskGroup で構造化並行性にヒント(段階的開示)
ヒント1 — 方向性
threading.Thread を asyncio に置き換える前に「なぜバックプレッシャーが必要か」を考える。プロデューサーが速くコンシューマーが遅いとき、無制限キューはメモリを食い潰す。asyncio.Queue(maxsize=N) は put() で自動ブロックするため、下流の処理速度に合わせて上流を制御できる(自然なバックプレッシャー)。contextvars.ContextVar は campaign_id を引数なしで全コルーチンに伝播できる仕組み。
ヒント2 — アプローチ
type JobItem = tuple[str, str](PEP 695)で型エイリアスを宣言的に定義_CAMPAIGN_CTX: ContextVar[str] = ContextVar("campaign_id", default="")でキャンペーン ID を暗黙伝播@dataclass(frozen=True, slots=True) class DispatchResult: sent_count: int; success_count: int; failed_ids: tuple[str, ...]- プロデューサー:
async def producer(queue: asyncio.Queue[JobItem], items: list[JobItem]) -> None - コンシューマー:
async def worker(worker_id: int, queue: asyncio.Queue[JobItem], results: WorkerResults) -> None async with asyncio.TaskGroup() as tg:で全タスクをまとめ、例外はExceptionGroupとして一括伝播- センチネル:
SENTINEL: Final[str] = "__DONE__"をワーカー数分投入して終了シグナル
ヒント3 — コードの骨格
import asyncio, contextvars
from dataclasses import dataclass
from typing import Final
# PEP 695 型エイリアス
type JobItem = tuple[str, str] # (recipient_id, message)
_CAMPAIGN_CTX: contextvars.ContextVar[str] = contextvars.ContextVar("campaign_id", default="")
MAX_QUEUE_SIZE: Final[int] = 100
NUM_WORKERS: Final[int] = 5
SENTINEL: Final[str] = "__DONE__"
@dataclass(frozen=True, slots=True)
class DispatchResult:
sent_count: int
success_count: int
failed_ids: tuple[str, ...]
@property
def success_rate(self) -> float:
return self.success_count / self.sent_count if self.sent_count else 0.0
async def producer(queue: asyncio.Queue[JobItem], items: list[JobItem]) -> None:
for item in items:
await queue.put(item) # maxsize 到達時は自動ブロック
for _ in range(NUM_WORKERS):
await queue.put((SENTINEL, "")) # sentinel
async def worker(wid: int, queue: asyncio.Queue[JobItem], results: list[DispatchResult]) -> None:
campaign_id = _CAMPAIGN_CTX.get()
...
async def run_campaign_async(campaign_id: str, jobs: list[JobItem]) -> DispatchResult:
token = _CAMPAIGN_CTX.set(campaign_id)
try:
queue: asyncio.Queue[JobItem] = asyncio.Queue(maxsize=MAX_QUEUE_SIZE)
partial_results: list[DispatchResult] = []
async with asyncio.TaskGroup() as tg:
tg.create_task(producer(queue, jobs))
for wid in range(NUM_WORKERS):
tg.create_task(worker(wid, queue, partial_results))
return aggregate_results(partial_results)
finally:
_CAMPAIGN_CTX.reset(token)
問題点分析(8点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | マジックナンバー散在 | 可読性 Ch7 | MAX_QUEUE_SIZE, NUM_WORKERS, SUCCESS_RATE_THRESHOLD 名前付き定数 |
| 2 | 型エイリアスなし | 型安全 Ch2 | type JobItem = tuple[str, str](PEP 695) |
| 3 | 返り値が dict | 型安全 Ch2/Ch4 | frozen dataclass DispatchResult |
| 4 | campaign_id の引数バケツリレー | 設計 Ch10 | contextvars.ContextVar で暗黙伝播 |
| 5 | マジックストリング "failed" | 型安全 Ch2 | StrEnum JobStatus |
| 6 | 集計が mutable dict + list | 不変性 Ch4 | frozen dataclass + tuple[str, ...] |
| 7 | queue.Queue() の maxsize 未指定 | バックプレッシャー Ch11 | asyncio.Queue(maxsize=MAX_QUEUE_SIZE) |
| 8 | threading.Thread — 例外伝播不可 | 構造化並行性 Ch10/Ch11 | asyncio.TaskGroup |
模範解答
import queue, threading, logging, time
# 問題1: マジックナンバー (100, 5, 0.05)
# 問題2: 型エイリアスなし
JOBS: list = []
def send_message(campaign_id, recipient_id, message):
# 問題3: 返り値が dict
time.sleep(0.01)
return {"status": "failed" if random.random() < 0.05 else "success"}
def worker_thread(job_queue, campaign_id, results, lock):
# 問題4: campaign_id の引数バケツリレー
failed, sent = [], 0
while True:
try: item = job_queue.get(timeout=1)
except: break
recipient_id, message = item
sent += 1
result = send_message(campaign_id, recipient_id, message)
if result["status"] == "failed": # 問題5: マジックストリング
failed.append(recipient_id)
job_queue.task_done()
with lock:
# 問題6: mutable dict + list
results.append({"sent": sent, "failed": failed})
def run_campaign(campaign_id, jobs):
# 問題7: maxsize 未指定 — バックプレッシャーなし
job_queue = queue.Queue()
for job in jobs: job_queue.put(job) # 全件一気に投入
results, lock, threads = [], threading.Lock(), []
for _ in range(5): # 問題8: threading.Thread — 例外伝播不可
t = threading.Thread(target=worker_thread, args=(job_queue, campaign_id, results, lock))
t.start(); threads.append(t)
job_queue.join()
for t in threads: t.join()
total_sent = sum(r["sent"] for r in results)
total_failed = sum(len(r["failed"]) for r in results)
return {"total_sent": total_sent, "total_success": total_sent - total_failed, "total_failed": total_failed}
"""campaign_dispatcher_async.py
Ch2: PEP 695 type JobItem / StrEnum JobStatus / 名前付き定数
Ch4: frozen dataclass slots=True DispatchResult(値オブジェクト)
Ch7: MAX_QUEUE_SIZE / NUM_WORKERS / SENTINEL 名前付き定数
Ch10: contextvars.ContextVar 暗黙伝播 / Fail-Fast
Ch11: asyncio.Queue Bounded Buffer / TaskGroup 構造化並行性 / reduce 集計
"""
from __future__ import annotations
import asyncio, contextvars, functools, logging, random
from dataclasses import dataclass
from enum import StrEnum
from typing import Final
logger = logging.getLogger(__name__)
# ── PEP 695 型エイリアス(Ch2)
type JobItem = tuple[str, str] # (recipient_id, message)
type WorkerResults = list[DispatchResult]
# ── 名前付き定数(Ch7)
MAX_QUEUE_SIZE: Final[int] = 100 # バックプレッシャー閾値
NUM_WORKERS: Final[int] = 5 # 並列ワーカー数
SUCCESS_RATE_THRESHOLD: Final[float] = 0.90
SENTINEL: Final[str] = "__DONE__" # ワーカー終了シグナル
# ── ContextVar(Ch10)
_CAMPAIGN_CTX: contextvars.ContextVar[str] = contextvars.ContextVar(
"campaign_id", default=""
)
# ── StrEnum(Ch2)
class JobStatus(StrEnum):
SUCCESS = "success"
FAILED = "failed"
SKIPPED = "skipped"
# ── 値オブジェクト(Ch4)
@dataclass(frozen=True, slots=True)
class DispatchResult:
"""配信ジョブの実行結果(イミュータブル値オブジェクト)。"""
sent_count: int
success_count: int
failed_ids: tuple[str, ...] # tuple: hashable + immutable
@property
def success_rate(self) -> float:
return self.success_count / self.sent_count if self.sent_count else 0.0
@property
def failure_count(self) -> int:
return self.sent_count - self.success_count
def is_healthy(self) -> bool:
return self.success_rate >= SUCCESS_RATE_THRESHOLD
# ── プロデューサー
async def producer(queue: asyncio.Queue[JobItem], items: list[JobItem]) -> None:
"""ジョブを asyncio.Queue に投入する。maxsize 到達時に自動ブロック。"""
campaign_id = _CAMPAIGN_CTX.get() # ContextVar から取得
logger.info("[Producer] campaign=%s total_jobs=%d", campaign_id, len(items))
for item in items:
await queue.put(item) # バックプレッシャー: 満杯時に自動ブロック
for _ in range(NUM_WORKERS):
await queue.put((SENTINEL, "")) # 終了シグナル
# ── ワーカー
async def worker(
worker_id: int,
queue: asyncio.Queue[JobItem],
results: WorkerResults,
) -> None:
"""キューからジョブを取り出して配信する非同期ワーカー。"""
campaign_id = _CAMPAIGN_CTX.get() # ContextVar から暗黙取得
sent, failed = 0, []
while True:
recipient_id, message = await queue.get()
if recipient_id == SENTINEL:
queue.task_done()
break
sent += 1
try:
await _send_stub(campaign_id, recipient_id, message)
except RuntimeError as exc:
logger.warning("[W%d] fail recipient=%s err=%s", worker_id, recipient_id, exc)
failed.append(recipient_id)
finally:
queue.task_done()
results.append(DispatchResult(
sent_count=sent,
success_count=sent - len(failed),
failed_ids=tuple(failed), # list → tuple でイミュータブル化
))
async def _send_stub(campaign_id: str, recipient_id: str, message: str) -> None:
"""外部送信 API の stub(95% 成功 / 5% 失敗)。"""
await asyncio.sleep(0)
if random.random() < 0.05:
raise RuntimeError(f"送信失敗: {recipient_id}")
# ── 集計(functools.reduce)
def aggregate_results(partial: WorkerResults) -> DispatchResult:
def merge(acc: DispatchResult, r: DispatchResult) -> DispatchResult:
return DispatchResult(
sent_count=acc.sent_count + r.sent_count,
success_count=acc.success_count + r.success_count,
failed_ids=acc.failed_ids + r.failed_ids,
)
return functools.reduce(merge, partial, DispatchResult(0, 0, ()))
# ── エントリポイント
async def run_campaign_async(campaign_id: str, jobs: list[JobItem]) -> DispatchResult:
"""非同期バックプレッシャー付きでキャンペーン配信を実行する。"""
token = _CAMPAIGN_CTX.set(campaign_id) # ContextVar をセット
try:
queue: asyncio.Queue[JobItem] = asyncio.Queue(maxsize=MAX_QUEUE_SIZE)
partial_results: WorkerResults = []
async with asyncio.TaskGroup() as tg: # 構造化並行性
tg.create_task(producer(queue, jobs))
for wid in range(NUM_WORKERS):
tg.create_task(worker(wid, queue, partial_results))
result = aggregate_results(partial_results)
if not result.is_healthy():
logger.warning("campaign=%s rate=%.1f%% below threshold", campaign_id, result.success_rate * 100)
return result
finally:
_CAMPAIGN_CTX.reset(token) # ContextVar を確実にリセット
import asyncio, logging, sys
logging.basicConfig(level=logging.INFO, stream=sys.stdout)
# 200 件のジョブを生成
sample_jobs: list[JobItem] = [
(f"user-{i:04d}", "夏セール開始!今すぐチェック CP-2026-SUMMER")
for i in range(200)
]
result = asyncio.run(run_campaign_async("CP-2026-SUMMER", sample_jobs))
print(f"sent={result.sent_count}")
print(f"success={result.success_count}")
print(f"failed={result.failure_count}")
print(f"rate={result.success_rate:.1%}")
print(f"healthy={result.is_healthy()}")
# INFO [Producer] campaign=CP-2026-SUMMER total_jobs=200
# INFO [W0] done campaign=CP-2026-SUMMER sent=40 failed=2
# INFO [W1] done campaign=CP-2026-SUMMER sent=40 failed=3
# ... (各ワーカー 40件ずつ処理)
# sent=200 success=190 failed=10 rate=95.0% healthy=True
# frozen → 変更不可
r = DispatchResult(sent_count=10, success_count=9, failed_ids=("u1",))
try:
r.sent_count = 999 # type: ignore
except Exception as e:
print(f"{type(e).__name__}: sent_count は変更不可") # FrozenInstanceError
# ContextVar はリセット確認
import contextvars
ctx: contextvars.ContextVar[str] = contextvars.ContextVar("test_ctx", default="")
token = ctx.set("CP-001")
print(f"set: {ctx.get()}") # CP-001
ctx.reset(token)
print(f"reset: {ctx.get()!r}") # ''
# バックプレッシャー確認: maxsize=3 の小さいキューで検証
async def backpressure_demo():
q: asyncio.Queue[tuple[str, str]] = asyncio.Queue(maxsize=3)
await q.put(("u1", "msg"))
await q.put(("u2", "msg"))
await q.put(("u3", "msg"))
print(f"qsize={q.qsize()} full={q.full()}") # qsize=3 full=True
# q.put(("u4", "msg")) をここで呼ぶと await するまでブロック
asyncio.run(backpressure_demo())
| ポイント | 適用した設計原則/パターン | 書籍対応章 |
|---|---|---|
type JobItem = tuple[str, str] PEP 695 型エイリアス | 型の活用・意図の明示 | Ch2 |
StrEnum JobStatus でマジックストリング排除 | 型安全・マジックストリング排除 | Ch2 |
@dataclass(frozen=True, slots=True) DispatchResult | 値オブジェクト・不変性・軽量化 | Ch4 |
tuple[str, ...] で failed_ids をイミュータブル化 | 不変性・型安全 | Ch4 |
MAX_QUEUE_SIZE, NUM_WORKERS, SENTINEL 名前付き定数 | マジックナンバー排除 | Ch7 |
ContextVar で campaign_id を暗黙伝播 | 関心の分離・引数バケツリレー排除 | Ch10 |
_CAMPAIGN_CTX.reset(token) で finally クリーンアップ | リソース管理・Fail-Safe | Ch10 |
asyncio.Queue(maxsize=N) Bounded Buffer | バックプレッシャー・メモリ安全 | Ch11 |
asyncio.TaskGroup 構造化並行性 | 例外伝播保証・タスクライフサイクル管理 | Ch11 |
functools.reduce(merge, partial, initial) 集計 | 関数型プログラミング・副作用ゼロ集計 | Ch11 |
# tests/test_campaign_dispatcher_async.py
import asyncio
import pytest
from campaign_dispatcher_async import (
JobItem, DispatchResult, JobStatus,
MAX_QUEUE_SIZE, NUM_WORKERS, SENTINEL, SUCCESS_RATE_THRESHOLD,
_CAMPAIGN_CTX, producer, worker, aggregate_results, run_campaign_async,
)
class TestDispatchResult:
def make(self, **kw):
defaults = dict(sent_count=10, success_count=9, failed_ids=("u1",))
return DispatchResult(**{**defaults, **kw})
def test_success_rate(self):
r = self.make(sent_count=10, success_count=8)
assert r.success_rate == pytest.approx(0.8)
def test_zero_sent(self):
r = self.make(sent_count=0, success_count=0, failed_ids=())
assert r.success_rate == 0.0
def test_frozen(self):
r = self.make()
with pytest.raises(Exception):
r.sent_count = 999 # type: ignore
def test_failed_ids_is_tuple(self):
r = self.make(failed_ids=("u1", "u2"))
assert isinstance(r.failed_ids, tuple)
def test_is_healthy_above_threshold(self):
r = self.make(sent_count=100, success_count=95, failed_ids=())
assert r.is_healthy()
def test_is_healthy_below_threshold(self):
r = self.make(sent_count=100, success_count=85, failed_ids=())
assert not r.is_healthy()
class TestAggregateResults:
def test_empty(self):
result = aggregate_results([])
assert result == DispatchResult(0, 0, ())
def test_multiple(self):
r1 = DispatchResult(100, 95, tuple(f"u{i}" for i in range(5)))
r2 = DispatchResult(100, 90, tuple(f"v{i}" for i in range(10)))
agg = aggregate_results([r1, r2])
assert agg.sent_count == 200
assert agg.success_count == 185
assert agg.failure_count == 15
class TestContextVar:
def test_context_propagation(self):
async def inner():
_CAMPAIGN_CTX.set("CP-TEST")
return _CAMPAIGN_CTX.get()
result = asyncio.run(inner())
assert result == "CP-TEST"
def test_context_reset(self):
async def inner():
token = _CAMPAIGN_CTX.set("CP-TEMP")
val1 = _CAMPAIGN_CTX.get()
_CAMPAIGN_CTX.reset(token)
val2 = _CAMPAIGN_CTX.get()
return val1, val2
v1, v2 = asyncio.run(inner())
assert v1 == "CP-TEMP"
assert v2 == "" # リセット後はデフォルト値
class TestBackpressure:
def test_queue_maxsize(self):
async def inner():
q: asyncio.Queue[JobItem] = asyncio.Queue(maxsize=3)
await q.put(("u1", "m"))
await q.put(("u2", "m"))
await q.put(("u3", "m"))
assert q.full()
assert q.qsize() == MAX_QUEUE_SIZE // (MAX_QUEUE_SIZE // 3)
asyncio.run(inner())
class TestRunCampaignAsync:
def test_full_run(self):
jobs: list[JobItem] = [(f"u{i}", "msg") for i in range(50)]
async def inner():
return await run_campaign_async("CP-TEST", jobs)
result = asyncio.run(inner())
assert result.sent_count == 50
assert result.success_count + result.failure_count == 50
assert isinstance(result.failed_ids, tuple)
設計図(Bounded Buffer / Producer-Consumer with TaskGroup)
ポイント解説
asyncio.Queue(maxsize=N)でバックプレッシャー(Ch11):maxsizeを設定するとqueue.put()はキューが満杯のとき自動的にawaitし、プロデューサーを下流ワーカーの処理速度に合わせて減速させる。メモリ無限増大を防ぐ唯一の簡潔な仕組みasyncio.TaskGroupで構造化並行性(Ch10/Ch11):async with TaskGroup() as tg:内のタスクが1つでも例外を投げると、残りのタスクを自動キャンセルしてExceptionGroupとして再送出する。gather(return_exceptions=True)と異なり、エラーが黙殺されないcontextvars.ContextVarで暗黙的コンテキスト伝播(Ch10):campaign_idを関数の引数として全コルーチンに渡すと引数が肥大化する。ContextVar.set()で設定しContextVar.get()で取得することで、ロギング・トレース文脈を自然に伝播できる。tokenをfinallyでリセットすることで再利用時の汚染を防ぐtype JobItem = tuple[str, str](PEP 695 / Ch2): Python 3.12+ のtype文は型エイリアスをTypeAliasより宣言的に記述できる。typeキーワードはスコープを持つため、モジュール内で再利用しやすく、mypy/pyright でも正しく推論される@dataclass(frozen=True, slots=True) DispatchResult(Ch4):frozen=TrueでFrozenInstanceErrorによる書き換え禁止、slots=Trueで__dict__を排除して軽量化。failed_ids: tuple[str, ...]はlistだとslots=Trueの hashable 要件と相性が悪いためtupleを選択
実務への応用
- ECサイト MOps キャンペーン配信: 100万件の宛先リストを
asyncio.Queueに流す際、maxsize=500程度でバックプレッシャーをかけ、外部メール送信 API(SendGrid 等)のレート制限と自然に同調させる - Argo Workflows: 各 Step の Python コンテナに本パターンを適用し、
DispatchResult.failed_idsを Argo のoutputs.parametersに渡してリドライブ(再送)ワークフローを起動できる - DataDog/OTel:
_CAMPAIGN_CTX.get()で取得したcampaign_idをspan.set_attribute("campaign.id", campaign_id)として付与し、キャンペーン別のスループット・エラー率をダッシュボードで可視化できる - GKE スケーリング:
NUM_WORKERSとMAX_QUEUE_SIZEを環境変数から PydanticBaseSettingsで注入し、Pod の CPU/メモリに合わせて動的チューニングする
今日のまとめ
asyncio.Queue(maxsize=N) によるバックプレッシャー + asyncio.TaskGroup による構造化並行性 + contextvars.ContextVar によるコンテキスト暗黙伝播の3点セットで、MOps キャンペーン配信バッチを「メモリ安全・例外伝播漏れなし・ロギング文脈統一」の非同期アーキテクチャに改善できる。