在 agent 系统中,经常需要一个组件产生数据,另一个组件消费数据。 比如: SSE 解析器 → 产出事件 → agentic loop 消费事件
()
| 413 | # ============================================================ |
| 414 | |
| 415 | def lesson_6_queue(): |
| 416 | """ |
| 417 | 在 agent 系统中,经常需要一个组件产生数据,另一个组件消费数据。 |
| 418 | 比如: SSE 解析器 → 产出事件 → agentic loop 消费事件 |
| 419 | """ |
| 420 | print("\n" + "=" * 60) |
| 421 | print("第六课: asyncio.Queue(生产者/消费者)") |
| 422 | print("=" * 60) |
| 423 | |
| 424 | async def sse_producer(queue: asyncio.Queue): |
| 425 | """生产者: 模拟 SSE 事件到达""" |
| 426 | events = ["message_start", "delta:Hello", "delta: World", "message_stop"] |
| 427 | for event in events: |
| 428 | await asyncio.sleep(0.1) |
| 429 | await queue.put(event) |
| 430 | print(f" [生产者] 放入: {event}") |
| 431 | await queue.put(None) # 结束信号 |
| 432 | |
| 433 | async def ui_consumer(queue: asyncio.Queue): |
| 434 | """消费者: 模拟 UI 逐字显示""" |
| 435 | text = "" |
| 436 | while True: |
| 437 | event = await queue.get() |
| 438 | if event is None: |
| 439 | break |
| 440 | if event.startswith("delta:"): |
| 441 | chunk = event[6:] |
| 442 | text += chunk |
| 443 | print(f" [消费者] 显示: '{chunk}' (累计: '{text}')") |
| 444 | else: |
| 445 | print(f" [消费者] 处理: {event}") |
| 446 | |
| 447 | async def demo(): |
| 448 | queue = asyncio.Queue() |
| 449 | # 生产者和消费者同时运行 |
| 450 | await asyncio.gather( |
| 451 | sse_producer(queue), |
| 452 | ui_consumer(queue), |
| 453 | ) |
| 454 | |
| 455 | print("\n 生产者(SSE)和消费者(UI)同时运行:") |
| 456 | asyncio.run(demo()) |
| 457 | |
| 458 | print(""" |
| 459 | 这个模式在 Claude Code 中的应用: |
| 460 | - SSE 解析器 → Queue → 消息构建器 |
| 461 | - Bash 命令输出 → Queue → 进度报告器 |
| 462 | - 多个 Worker agent → Queue → Leader agent 收集结果 |
| 463 | """) |
| 464 | |
| 465 | |
| 466 | # ============================================================ |
no test coverage detected