A コーディング — asyncio.Queue Bounded Buffer × asyncio.TaskGroup 構造化並行性 × contextvars.ContextVar 暗黙伝播 × PEP 695 type 文 × frozen dataclass slots=True DispatchResult(キャンペーン配信バッチ バックプレッシャー Bad→Good Ch2/Ch4/Ch7/Ch10/Ch11)

2026-06-29 (Day 84) 月曜 コーディング ★★★★☆ Python 3.12 / asyncio / contextvars / PEP 695 / frozen dataclass 良いコード設計入門 Ch2 / Ch4 / Ch7 / Ch10 / Ch11

概要

🚦

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() で取得することで、ロギング・トレース文脈を自然に伝播。tokenfinally でリセットし汚染を防ぐ。

🏷️

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.TaskGroupcontextvars.ContextVartype 文(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.ContextVarcampaign_id をコルーチン間の暗黙的コンテキストとして伝播すること
  • type JobItem = tuple[str, str] の PEP 695 型エイリアスを使うこと
  • 送信結果は @dataclass(frozen=True, slots=True)DispatchResult で返すこと
  • Google スタイル docstring・インラインコメント・名前付き定数を含めること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード(Google スタイル docstring・インラインコメント・名前付き定数含む)+ 実行例(input→output)+ 適用した設計パターン名と書籍対応章

悪いコード (Before)

このコードには 8つの設計上の問題 が隠れています。見つけてみてください。
bad_campaign_dispatcher.py — threading.Thread + 無制限キュー + 型なし + campaign_id 引数バケツリレー
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,
    }
問題点サマリー(8点)
1マジックナンバー散在(Ch7)100, 5, 0.05 がコードに埋め込まれている。MAX_QUEUE_SIZE, NUM_WORKERS, SUCCESS_RATE_THRESHOLD で名前付き定数に
2型エイリアスなし(Ch2)list が何のリストか不明。type JobItem = tuple[str, str] の PEP 695 型エイリアスで意図を明示
3返り値が dict(Ch2/Ch4) — キーのタイポを型チェッカーが検出できない。frozen dataclass DispatchResult で型安全な値オブジェクトに
4campaign_id の引数バケツリレー(Ch10) — 全関数シグネチャが肥大化。contextvars.ContextVar で暗黙的に伝播
5マジックストリング "failed"(Ch2)StrEnum JobStatus で型安全に管理
6集計が mutable dict + list(Ch4)frozen=True, slots=True DispatchResult + tuple[str, ...] でイミュータブル化
7queue.Queue() の maxsize 未指定(Ch11) — バックプレッシャーなしでメモリ無限増大。asyncio.Queue(maxsize=MAX_QUEUE_SIZE) で制御
8threading.Thread — 例外伝播不可(Ch10/Ch11) — スレッドの例外は親スレッドに伝わらず、キャンセルも困難。asyncio.TaskGroup で構造化並行性に

ヒント(段階的開示)

ヒント1 — 方向性
threading.Threadasyncio に置き換える前に「なぜバックプレッシャーが必要か」を考える。プロデューサーが速くコンシューマーが遅いとき、無制限キューはメモリを食い潰す。asyncio.Queue(maxsize=N)put() で自動ブロックするため、下流の処理速度に合わせて上流を制御できる(自然なバックプレッシャー)。contextvars.ContextVarcampaign_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マジックナンバー散在可読性 Ch7MAX_QUEUE_SIZE, NUM_WORKERS, SUCCESS_RATE_THRESHOLD 名前付き定数
2型エイリアスなし型安全 Ch2type JobItem = tuple[str, str](PEP 695)
3返り値が dict型安全 Ch2/Ch4frozen dataclass DispatchResult
4campaign_id の引数バケツリレー設計 Ch10contextvars.ContextVar で暗黙伝播
5マジックストリング "failed"型安全 Ch2StrEnum JobStatus
6集計が mutable dict + list不変性 Ch4frozen dataclass + tuple[str, ...]
7queue.Queue() の maxsize 未指定バックプレッシャー Ch11asyncio.Queue(maxsize=MAX_QUEUE_SIZE)
8threading.Thread — 例外伝播不可構造化並行性 Ch10/Ch11asyncio.TaskGroup

模範解答

Before — threading.Thread + 無制限キュー + 型なし + 引数バケツリレー
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}
After — asyncio.Queue × TaskGroup × ContextVar × PEP 695 × frozen dataclass
"""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-SafeCh10
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)

_CAMPAIGN_CTX: ContextVar[str] campaign_id を全コルーチンに暗黙伝播 producer() ジョブをキューに投入 await queue.put() でブロック asyncio.Queue(maxsize=100) Bounded Buffer — バックプレッシャー制御 put() → 満杯時に自動ブロック asyncio.TaskGroup worker(0) — await queue.get() worker(1) — await queue.get() worker(2) — await queue.get() worker(3) — await queue.get() worker(4) — await queue.get() 例外発生 → 全タスクを自動キャンセル ContextVar.get() @dataclass(frozen=True, slots=True) DispatchResult sent_count, success_count, failed_ids: tuple[str, ...] success_rate (property) / is_healthy() / failure_count (property) results.append() aggregate_results() functools.reduce tuple 連結 + カウント合算 Sentinel ("__DONE__") を NUM_WORKERS 個投入してワーカーを終了

ポイント解説

  1. asyncio.Queue(maxsize=N) でバックプレッシャー(Ch11): maxsize を設定すると queue.put() はキューが満杯のとき自動的に await し、プロデューサーを下流ワーカーの処理速度に合わせて減速させる。メモリ無限増大を防ぐ唯一の簡潔な仕組み
  2. asyncio.TaskGroup で構造化並行性(Ch10/Ch11): async with TaskGroup() as tg: 内のタスクが1つでも例外を投げると、残りのタスクを自動キャンセルして ExceptionGroup として再送出する。gather(return_exceptions=True) と異なり、エラーが黙殺されない
  3. contextvars.ContextVar で暗黙的コンテキスト伝播(Ch10): campaign_id を関数の引数として全コルーチンに渡すと引数が肥大化する。ContextVar.set() で設定し ContextVar.get() で取得することで、ロギング・トレース文脈を自然に伝播できる。tokenfinally でリセットすることで再利用時の汚染を防ぐ
  4. type JobItem = tuple[str, str](PEP 695 / Ch2): Python 3.12+ の type 文は型エイリアスを TypeAlias より宣言的に記述できる。type キーワードはスコープを持つため、モジュール内で再利用しやすく、mypy/pyright でも正しく推論される
  5. @dataclass(frozen=True, slots=True) DispatchResult(Ch4): frozen=TrueFrozenInstanceError による書き換え禁止、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_idspan.set_attribute("campaign.id", campaign_id) として付与し、キャンペーン別のスループット・エラー率をダッシュボードで可視化できる
  • GKE スケーリング: NUM_WORKERSMAX_QUEUE_SIZE を環境変数から Pydantic BaseSettings で注入し、Pod の CPU/メモリに合わせて動的チューニングする

今日のまとめ

asyncio.Queue(maxsize=N) によるバックプレッシャー + asyncio.TaskGroup による構造化並行性 + contextvars.ContextVar によるコンテキスト暗黙伝播の3点セットで、MOps キャンペーン配信バッチを「メモリ安全・例外伝播漏れなし・ロギング文脈統一」の非同期アーキテクチャに改善できる。

自己評価(あとで記入)

理解度

自分の回答

気づき・メモ