跳转至

ADR-0002:采用 Redis Stream 作为消息总线

状态

已接受

背景

代理商 backend 需要异步处理短信发送、向浏览器推送 SSE 事件。同步调用会阻塞 API,且不利于水平扩展 Worker。

决策

src/messaging/ 引入独立 Python 包,使用 Redis StreamXADD / XREADGROUP / XACK)实现发布订阅:

  1. 统一信封 StreamEnvelopeevent_type + payload + correlation_id
  2. 业务 Streamstream:sms(短信)、stream:sse(SSE 广播)
  3. 消费者组 messaging-workers,由 backend 进程 lifespan 或独立 Worker 消费
  4. 测试/无 RedisMESSAGING_USE_MEMORY=true 使用 MemoryStreamBackend
  5. 关闭总线MESSAGING_ENABLED=falsebackend 直调 send_code

架构约束

  • messaging 不依赖 backend / myapp
  • backend.core 不直接 import messaging;wiring 放在 backend/composition/messaging.py

后果

优点

  • API 与短信/SSE 推送解耦,可独立扩 Worker
  • Stream 持久化与 consumer group 支持至少一次投递
  • 内存后端便于单元测试

缺点

  • 生产环境需部署 Redis
  • SSE 经 Stream 中转,延迟略高于进程内广播(可接受)

备选方案

方案 未采用原因
Redis Pub/Sub 无持久化,消费者离线丢消息
Celery/RQ 引入较重,当前队列场景较简单
进程内 asyncio.Queue 无法跨进程/跨实例