跳转至

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 MessagingSettingsMESSAGING_* / REDIS_URL
messaging/schemas/ StreamEnvelope 统一消息信封
messaging/streams/ 发布器、消费者、Redis/内存后端
messaging/workers/ SMS 处理器、SSE Hub、后台 Runner
messaging/composition/ 进程内单例与 FastAPI Depends

依赖规则

  • messaging 不得 import backendmyapp(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

相关文档