Redis Stream 消息组件¶
src/messaging/ 是基于 Redis Stream 的异步发布订阅组件,供 backend 等业务包使用,典型场景包括:
- 短信任务队列:API 发布
sms.send_code,Worker 消费并调用backend.utils.sms(支持mock/ 阿里云号码认证dypns) - 大模型生图队列:
POST /api/canvas/generate发布llm.image_generate,Worker 调用 NoneLinear API - SSE 推送:业务事件写入
stream:sse,客户端通过GET /api/events/stream长连接接收
目录结构¶
| 路径 | 用途 |
|---|---|
messaging/config.py |
MessagingSettings(MESSAGING_* / REDIS_URL) |
messaging/schemas/ |
StreamEnvelope 统一消息信封 |
messaging/streams/ |
发布器、消费者、Redis/内存后端 |
messaging/workers/ |
SMS 处理器、SSE Hub、后台 Runner |
messaging/composition/ |
进程内单例与 FastAPI Depends |
依赖规则¶
messaging不得 importbackend或myapp(import-linter 强制)backend通过backend/composition/messaging.py完成 wiring
环境变量¶
| 变量 | 说明 | 默认 |
|---|---|---|
MESSAGING_ENABLED |
是否启用 Stream(false 时 API 直调短信) |
true |
MESSAGING_USE_MEMORY |
无 Redis 时使用内存 Stream | false |
REDIS_URL / MESSAGING_REDIS_URL |
Redis 连接 | redis://127.0.0.1:6379/0 |
MESSAGING_CONSUMER_GROUP |
消费者组 | messaging-workers |
MESSAGING_SMS_STREAM |
短信 Stream 键 | stream:sms |
MESSAGING_LLM_STREAM |
大模型/生图 Stream 键 | stream:llm |
SMS_PROVIDER |
短信模式:mock / dypns(号码认证) |
mock |
SMS_SIGN_NAME |
Dypns 赠送签名 | — |
SMS_TEMPLATE_CODE |
Dypns 赠送模板 Code | — |
SMS_ACCESS_KEY / SMS_SECRET_KEY |
RAM 凭证(需 SendSmsVerifyCode / CheckSmsVerifyCode) |
— |
MESSAGING_SSE_STREAM |
SSE Stream 键 | stream:sse |
单元测试默认 MESSAGING_ENABLED=false,避免依赖 Redis。
数据流¶
sequenceDiagram
participant API as backend API
participant Pub as StreamPublisher
participant Redis as Redis Stream
participant Worker as SMS Worker
participant SMS as sms.send_code
API->>Pub: publish(sms.send_code)
Pub->>Redis: XADD stream:sms
Worker->>Redis: XREADGROUP
Worker->>SMS: send_code(phone)
Worker->>Redis: XACK
SSE¶
- 发布:
SseHub.publish(event, data, user_id=...) - 订阅:
GET /api/events/stream(需登录,按user_id过滤)
本地开发¶
# 启动 Redis(示例)
docker run -d --name redis -p 6379:6379 redis:7
# 终端 1:API
make run-backend
# 终端 2:Stream Worker(短信、生图队列消费)
make run-backend-worker
无 Redis 时可设置:
MESSAGING_ENABLED=true
MESSAGING_USE_MEMORY=true
Kubernetes¶
| 组件 | 说明 |
|---|---|
redis StatefulSet |
Stream 持久化后端 |
backend-worker Deployment |
Init 等待 Redis;主容器 python -m backend.worker_main |
MESSAGING_ENABLED |
ConfigMap 默认 true |
排错:
kubectl get pods -n harness-app -l 'app in (redis,backend-worker)'
kubectl logs -n harness-app deploy/backend-worker -c worker
ACK 上 Redis PVC 需 ESSD ≥ 20Gi,见 ACK Runbook — Redis PVC。