# 后端消息队列改造设计 **日期:** 2026-07-01 **状态:** 设计完成,待实施 **分支:** feat/actvity --- ## 一、背景与目标 ### 1.1 现状问题 | 问题 | 具体表现 | 影响 | |------|----------|------| | 同步 RPC 紧耦合 | socialService 点赞后同步调 notificationService,notificationService 挂了影响点赞 | 故障传播 | | fire-and-forget 裸 goroutine | notificationService.push + moderationService 举报通知均用 `go func()` 发推送/通知,无重试无持久化 | 推送丢失、通知丢失 | | 内存 channel | statisticService 用 Go channel 传事件,服务重启清空 | 埋点数据丢失 | | 轮询替代事件驱动 | galleryService 每分钟 ticker 查过期展览 | 收益结算延迟 | | 手动下架不计累计时长 | exhibitionService.RemoveFromSlot 只调 RemoveExhibitionTx,**不调** userService.AddExhibitionHours | 主动下架的展示时长、收益、点赞押注全部丢失 | | cron worker 自管理 | taskService/assetService 各自 sleep loop 实现定时任务 | 不可靠、难监控 | ### 1.2 目标 1. **解耦** — 跨服务异步通信替换同步 RPC 2. **可靠** — 所有异步操作有重试/死信/持久化 3. **实时** — 事件驱动替代轮询 4. **不引入新中间件** — 基于现有 Redis --- ## 二、技术选型 ### 2.1 Asynq — 任务队列 - 选型理由:基于 Redis(已有),Go 原生,支持重试/超时/死信/延时/定时 - 适用场景:点对点任务,需要重试保证 ### 2.2 Redis Streams — 事件流 - 选型理由:基于 Redis(已有),activityService 已有使用经验,支持消费者组 ACK - 适用场景:一对多事件,多消费者、高吞吐 ### 2.3 决策原则 | 场景 | 用 | |------|-----| | 点对点、需要重试/死信 | Asynq | | 一对多、多消费者、高吞吐 | Redis Streams | | 读操作、强一致性写 | 保留同步 RPC | --- ## 三、架构与调用关系 ### 3.1 改造前(现状) ``` ┌─────────────────────────────────────────────────────────────────────────┐ │ Gateway (Gin HTTP) │ │ Dubbo 客户端,全部走同步 RPC 调用下游 │ └────┬────┬────┬────┬────┬────┬────┬────┬────┬────┬────┬────┬────┬────────┘ │ │ │ │ │ │ │ │ │ │ │ │ │ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────────┐ ┌────────┐ ┌─────────┐ ┌──────────┐ ┌────────┐ ┌──────────┐ ┌──────────┐ │ user │ │ asset │ │gallery │ │ task │ │notification│ │ social │ │moderation│ │statistic │ │starbook│ │aiChat │ │activity │ │Service │ │Service │ │Service │ │Service │ │ Service │ │Service │ │ Service │ │ Service │ │Service │ │Service │ │Service │ └───┬────┘ └───┬────┘ └───┬────┘ └───┬────┘ └─────┬──────┘ └───┬────┘ └────┬─────┘ └────┬─────┘ └────────┘ └────────┘ └────┬─────┘ │ │ │ │ │ │ │ │ │ │ │ cron │ ticker │ sleep │ go func() │ │ │ channel │ │ │ worker │ 1 min │ loop │ push │ │ │ sink │ │ │ │ │ │ │ │ │ │ └──────────┴──────────┴──────────┴─────────────┴────────────┴──────────┴───────────┴──────────────────────────────┘ 全部通过 Dubbo 同步 RPC 直连,失败只打日志,无重试 ``` **核心问题标注:** | 标注 | 位置 | 问题 | |------|------|------| | `cron worker` | assetService | `season_reset_worker` 只靠 cron 定时 | | `ticker 1 min` | galleryService | `CleanupWorker` 每分钟轮询过期展览 → 同步 RPC taskService | | `sleep loop` | taskService | `DailyResetWorker` goroutine 计算到 05:00 的等待时间 | | `go func() push` | notificationService | 废弃 goroutine 发推送,无重试 | | `channel sink` | statisticService | 内存 `chan *Event` → 重启丢数据 | ### 3.2 改造后(目标架构) ``` ┌─────────────────────────┐ │ Gateway (Gin) │ │ ┌───────────────────┐ │ │ │ WebSocket Hub │◄──┼── Redis Streams │ │ (stream:activity) │ │ Consumer Group │ └───────────────────┘ │ └────┬────────────────────┘ │ Dubbo RPC(读/强一致性写保持同步) │ ┌─────────────────────────┼─────────────────────────┐ │ │ │ ▼ ▼ ▼ ┌─────────────────┐ ┌──────────────────┐ ┌──────────────────┐ │ Asynq Broker │ │ Redis Streams │ │ Dubbo RPC │ │ (Redis DB 2) │ │ (Redis DB 3) │ │ (读/同步写) │ │ │ │ │ │ │ │ ▪ 任务队列 │ │ ▪ 事件流 │ │ userService │ │ ▪ 重试+死信 │ │ ▪ 消费者组 ACK │ │ assetService │ │ ▪ 延时/定时 │ │ ▪ 多消费者 │ │ galleryService │ │ ▪ 优先级 │ │ ▪ 消息持久化 │ │ socialService │ └───┬─────────────┘ └───┬───────────────┘ │ ... 等读接口 │ │ │ └──────────────────┘ │ │ │ ┌──────────────────────────────────────────────────────────┐ │ │ Asynq 生产者 / 消费者 │ ├──┤ │ │ │ 生产者: │ │ │ galleryService → notification:create │ │ │ → revenue:exhibition │ │ │ → revenue:like-bet │ │ │ → gallery:exhibition-expire (delay) │ │ │ socialService → notification:create │ │ │ moderationService → notification:create │ │ │ → moderation:auto-hide │ │ │ aiChatService → aichat:chat │ │ │ starbookService → starbook:collection-create │ │ │ │ │ │ 消费者 + Scheduler: │ │ │ notificationService ← notification:create │ │ │ ← notification:push (self) │ │ │ taskService ← revenue:exhibition │ │ │ ← revenue:like-bet │ │ │ ← task:daily-reset (scheduler) │ │ │ galleryService ← gallery:exhibition-expire (self) │ │ │ assetService ← asset:season-reset (scheduler) │ │ │ statisticService ← statistic:materialize (scheduler) │ │ │ ← statistic:weekly-income (scheduler)│ │ │ ← statistic:level-up (scheduler) │ │ │ moderationService ← moderation:auto-hide (self) │ │ │ aiChatService ← aichat:chat (self) │ │ │ starbookService ← starbook:collection-create (self) │ │ └──────────────────────────────────────────────────────────┘ │ │ ┌──────────────────────────────────────────────────────────┐ │ │ Redis Streams 生产者 / 消费者 │ ├──┤ │ │ │ 生产者 → XADD: │ │ │ userService → stream:user │ │ │ assetService → stream:asset │ │ │ socialService → stream:social │ │ │ galleryService → stream:exhibition │ │ │ activityService→ stream:activity │ │ │ moderationService → stream:moderation │ │ │ │ │ │ 消费者组 (Consumer Group: topfans-service): │ │ │ statisticService ← stream:user │ │ │ ← stream:asset │ │ │ ← stream:social │ │ │ ← stream:exhibition │ │ │ ← stream:moderation │ │ │ taskService ← stream:user (用户注册 → 初始化任务) │ │ │ ← stream:asset (铸造统计) │ │ │ ← stream:exhibition (展览统计) │ │ │ notificationService ← stream:moderation (举报通知) │ │ └──────────────────────────────────────────────────────────┘ ``` ### 3.3 关键调用链 **展示收益结算链(改造后):** ``` 展览上架 galleryService │ ├─ 写 ZSET(保留兼容) └─ mq.EnqueueExhibitionExpire(exhibition_id, expireAt) ← 走 adapter │ ▼ ┌─ 到期时刻触发 ─────────────────────────────┐ │ galleryService (handler) │ │ ├─ 查询点赞数 │ │ ├─ adapter.Get().EventProducer() │ │ │ .Publish("stream:exhibition", e) │ │ ├─ mq.EnqueueRevenueExhibition(payload) │ │ ├─ mq.EnqueueRevenueLikeBet(payload) │ │ ├─ 标记 exhibition.processed=true │ │ └─ ZSET.Remove │ │ │ ▼ ┌─ revenue:exhibition ────────────────────────┐ │ taskService (handler) │ │ ├─ assetLevelService.CalculateRevenue() RPC │ │ ├─ CreateRevenueRecord() → DB │ │ ├─ assetLevelService.AddExhibitionHours() RPC│ │ ├─ userService.AddExhibitionHours() RPC │ │ └─ 失败 → 重试3次 → 死信队列 │ │ │ ▼ ┌─ revenue:like-bet ──────────────────────────┐ │ taskService (handler) │ │ ├─ 查询 exhibition 下所有 asset_likes │ │ ├─ 按 bet_order 计算每笔金额 │ │ ├─ BatchCreate like_bet_revenue_records │ │ └─ DB 唯一约束保证幂等 │ └────────────────────────────────────────────────┘ ``` **通知推送链(改造后):** ``` 点赞操作 socialService ├─ assetClient.LikeAsset() RPC ✅ 同步 └─ mq.EnqueueLikeNotification(...) ← 走 adapter │ ▼ notificationService (handler) ├─ 事务内 INSERT notification + UPSERT stats └─ mq.EnqueuePushPayload(notif) ← 走 adapter │ ▼ notificationService (handler) ├─ 拉取用户活跃 cids ├─ UniPush.Send() ├─ 失败 → 重试5次(指数退避) └─ 最终失败 → 死信队列 → 告警 ``` ### 3.4 服务依赖关系(改造后) ``` ┌──────────┐ │ Redis │ ← 唯一的中间件依赖 └────┬─────┘ ┌─────────────┼─────────────┐ │ │ │ ┌────▼────┐ ┌─────▼──────┐ ┌──▼──────────┐ │ Asynq │ │ Streams │ │ Dubbo RPC │ │ (DB 2) │ │ (DB 3) │ │ (同步读/写) │ └────┬────┘ └─────┬──────┘ └──┬──────────┘ │ │ │ ┌──────▼──────┐ │ ┌──────▼──────┐ │ 任务型服务 │ │ │ 查询型服务 │ │ notif/task/ │ │ │ user/asset/ │ │ gallery等 │ │ │ social等 │ └─────────────┘ │ └─────────────┘ ┌──────▼──────┐ │ 事件消费型 │ │ statistic/ │ │ gateway(WS) │ └─────────────┘ ``` - **Asynq**:notificationService、taskService、galleryService 依赖(既是生产者也是消费者) - **Redis Streams**:statisticService 重度依赖(消费者组),userService/assetService/socialService 轻量依赖(仅生产) - **Dubbo RPC**:保留用于读操作和强一致性写(如领取收益时调 userService.UpdateCrystalBalance) - **不新增中间件**:Asynq 和 Streams 都跑在现有 Redis 上 --- ## 四、新增包结构(适配器模式 — broker 可插拔) ### 4.0 设计原则 > **业务代码只跟 `pkg/mq/adapter/` 接口层打交道,不直接依赖 Asynq / Redis Streams / RabbitMQ 的具体类型。** > 后续切换 broker(如从 Redis 迁到 RabbitMQ)只需要新增一个 adapter 实现,业务侧 producer/consumer 一行不改。 ``` backend/pkg/mq/ # 共享:broker 无关的通用契约 ├── adapter/ # ★ 适配器抽象层 — 业务代码面向这些接口编程 │ ├── adapter.go # Adapter 接口 + Manager 单例 │ ├── task.go # Task / TaskInfo 公共数据结构 │ ├── event.go # Event / MessageAck 公共数据结构 │ ├── task_producer.go # TaskProducer 接口:Enqueue / EnqueueAt │ ├── task_consumer.go # TaskConsumer 接口:RegisterTask / Run │ ├── event_producer.go # EventProducer 接口:Publish │ ├── event_consumer.go # EventConsumer 接口:Subscribe / Ack │ └── options.go # 通用配置选项(重试次数、退避、延迟等) │ ├── asynq/ # Asynq adapter 实现(当前默认) │ ├── client.go # Asynq Client 单例 │ ├── server.go # Asynq Server 单例 │ ├── task_producer.go # 实现 adapter.TaskProducer │ ├── task_consumer.go # 实现 adapter.TaskConsumer │ ├── task_marshal.go # Task payload 序列化/反序列化(JSON) │ └── middleware.go # 日志/重试中间件 │ ├── streams/ # Redis Streams adapter 实现 │ ├── event_producer.go # 实现 adapter.EventProducer │ ├── event_consumer.go # 实现 adapter.EventConsumer │ ├── keys.go # 所有 Stream Key 常量 │ └── consumer_group.go # 消费者组管理 │ ├── rabbitmq/ # ★ 预留(后续接入用,本期不实现) │ ├── task_producer.go # 实现 adapter.TaskProducer(占位) │ ├── task_consumer.go # 实现 adapter.TaskConsumer(占位) │ ├── event_producer.go # 实现 adapter.EventProducer(占位) │ └── event_consumer.go # 实现 adapter.EventConsumer(占位) │ ├── tasks/ # 业务侧 Task Type 注册中心(broker 无关) │ └── registry.go # 所有 TaskType 常量、Payload 结构定义 │ └── config.go # MQ 配置 + adapter 选择 backend/services/xxxService/ └── mq/ # 各服务 MQ 适配:业务侧 ├── producer.go # 通过 adapter.TaskProducer / EventProducer 发送 └── consumer.go # 通过 adapter.TaskConsumer / EventConsumer 注册 handler ``` ### 4.1 Adapter 接口设计(关键代码骨架) **adapter/adapter.go — 抽象工厂 + 选择器** ```go package adapter import "context" // Adapter 顶层抽象 — 一个进程内只有一个 adapter 实例 type Adapter interface { // 业务模型划分:两类原语 // - Task:点对点 + 重试 + 死信(→ Asynq / RabbitMQ-RabbitMQ) // - Event:发布订阅 + 多消费者(→ Redis Streams / RabbitMQ-Topic) TaskProducer() TaskProducer TaskConsumer() TaskConsumer EventProducer() EventProducer EventConsumer() EventConsumer Close() error } // Manager 单例,启动时根据 config 选择 adapter 实现 type Manager struct { impl Adapter } func Init(cfg Config) error { /* 选择 adapter 实现 */ } func Get() Adapter { return mgr.impl } // broker 选择由 config.MQDriver 决定 // 未来切到 RabbitMQ:mgr.impl = rabbitmq.NewAdapter(cfg) // 业务代码 0 改动 ``` **adapter/task.go — 业务侧 Task 数据结构** ```go package adapter import "time" // Task 业务侧定义(屏蔽 Asynq 的 asynq.Task) // 切 broker 时业务代码不感知 type Task struct { Type string // 任务类型 e.g. "notification:create" Payload map[string]any // JSON 业务参数 Queue string // 可选,默认为 "default" MaxRetry int // 重试上限 Timeout time.Duration // 单次执行超时 Delay time.Duration // 延时(ProcessAt 也能设) Priority int // 优先级 } type TaskInfo struct { ID string Queue string State string // pending / active / completed / failed NextRunAt time.Time Retry int } ``` **adapter/task_producer.go — 生产端接口** ```go package adapter import "context" type TaskProducer interface { // Enqueue 立即入队 Enqueue(ctx context.Context, t Task) (string, error) // EnqueueAt 延时入队,ProcessAt 时刻才执行(替代 ticker) EnqueueAt(ctx context.Context, t Task, processAt time.Time) (string, error) // EnqueueUnique 仅一个未执行任务存在(按 Type+UniqueKey 去重) EnqueueUnique(ctx context.Context, t Task, uniqueTTL time.Duration) (string, error) } ``` **adapter/task_consumer.go — 消费端接口** ```go package adapter type TaskHandler func(ctx context.Context, t *Task) error type TaskConsumer interface { RegisterTask(taskType string, handler TaskHandler, opts TaskRegisterOptions) error RegisterCron(spec string, taskType string, payload map[string]any) error // cron 定时 Run(ctx context.Context) error // 启动 worker Stop() error } type TaskRegisterOptions struct { MaxRetry int Queue string Timeout time.Duration } ``` **adapter/event.go — 事件数据结构** ```go package adapter // Event 业务侧事件结构(屏蔽 Redis Streams 的 XMessage) // 切 broker 时业务代码不感知 type Event struct { Type string // 事件类型 e.g. "asset.mint" Source string // 来源服务 OccurredAt time.Time Payload map[string]any // JSON } ``` **adapter/event_producer.go + event_consumer.go** ```go type EventProducer interface { Publish(ctx context.Context, topic string, e Event) error } type EventConsumer interface { Subscribe(ctx context.Context, topics []string, group string, handler EventHandler) error Ack(ctx context.Context, topic string, msgID string) error } type EventHandler func(ctx context.Context, topic string, e Event, msgID string) error ``` ### 4.2 各服务的 `mq/` 目录(业务侧,broker 无关) 每个服务只 import `pkg/mq/adapter`,**不直接 import asynq/streams 包**。 ``` backend/services/galleryService/mq/ ├── producer.go │ import "github.com/topfans/backend/pkg/mq/adapter" │ │ func EnqueueExhibitSettled(ctx context.Context, e ExhibitionEvent) error { │ return adapter.Get().TaskProducer().EnqueueAt( │ ctx, │ adapter.Task{ │ Type: "gallery:exhibition-settled", │ Payload: map[string]any{...e...}, │ MaxRetry: 3, │ }, │ e.SettledAt, │ ) │ } │ └── consumer.go import "github.com/topfans/backend/pkg/mq/adapter" func RegisterHandlers() error { consumer := adapter.Get().TaskConsumer() consumer.RegisterTask("gallery:exhibition-settled", handleSettled, adapter.TaskRegisterOptions{MaxRetry: 3}) consumer.RegisterTask("gallery:exhibition-expire", handleExpire, adapter.TaskRegisterOptions{MaxRetry: 3}) consumer.RegisterCron("0 */1 * * *", "gallery:cleanup-display-status", nil) return nil } ``` ### 4.3 Asynq Adapter 实现关键点(asynq/ 目录) | 文件 | 职责 | |------|------| | `client.go` | `asynq.Client` 单例,封装 Redis 连接配置 | | `server.go` | `asynq.Server` 单例,封装 concurrency / queues | | `task_producer.go` | `adapter.Task` → `asynq.Task`,调 `asynq.Client.Enqueue`;`EnqueueAt` 用 `ProcessAt` 实现 | | `task_consumer.go` | 业务侧 `TaskHandler` → `asynq.HandlerFunc`,注册到 `asynq.Mux` | | `task_marshal.go` | `adapter.Task.Payload` (map) ↔ JSON 序列化,确保反序列化两端兼容 | ### 4.4 Redis Streams Adapter 实现关键点(streams/ 目录) | 文件 | 职责 | |------|------| | `event_producer.go` | `adapter.Event` → `XADD stream:user ...`,Payload map → fields | | `event_consumer.go` | 包装 `XReadGroup`,收到消息后回调业务 `EventHandler` | | `keys.go` | 所有 stream key 常量(`stream:user` 等) | | `consumer_group.go` | `XGroupCreateMkStream`,自动重试,pending list 清理 | ### 4.5 RabbitMQ 接入路径(本期不实现,仅留位) 未来切到 RabbitMQ 时: 1. 在 `pkg/mq/rabbitmq/` 下实现 4 个文件 2. `pkg/mq/config.go` 改 `Driver` = `"rabbitmq"` 3. adapter.Manager 自动切换 4. 业务代码(所有服务的 `mq/producer.go` 和 `mq/consumer.go`)**0 改动** ### 4.6 完整调用链示意(接入透明) ``` 业务方: mq.PublishExhibitSettled(ctx, exhibitionData) ↓ adapter.Get().TaskProducer().EnqueueAt(...) ← 业务只跟 adapter 打交道 ↓ [manager impl = asynq] asynq.TaskProducer.EnqueueAt() ↓ asynq.Client.Enqueue() ← 真正的 MQ 调用 切换到 RabbitMQ: config.MQDriver = "rabbitmq" ↓ [manager impl = rabbitmq] rabbitmq.TaskProducer.EnqueueAt() ↓ amqp091-go.Channel.Publish(...) 业务方代码:不变 ``` --- ## 五、Redis Streams 定义 | Stream Key | 生产者 | 消费者 | 事件描述 | |------------|--------|--------|----------| | `stream:user` | userService | statisticService, taskService | 注册、资料变更 | | `stream:asset` | assetService | statisticService, taskService | 铸造、等级变更 | | `stream:social` | socialService | statisticService | 资产点赞 | | `stream:exhibition` | galleryService | taskService, statisticService | 上架开始/到期/完成 | | `stream:activity` | activityService | gateway(WebSocket) | 活动贡献(已有,统一命名) | | `stream:moderation` | moderationService | notificationService | 举报/处理结果 | > `stream:activity` 对应 activityService 已有的 `combo:stream:contributions`,逻辑不动,统一 key 命名。 --- ## 六、Asynq Task 定义 ### 6.1 通知类 | Task Type | 生产者 | 消费者 | 重试 | 说明 | |-----------|--------|--------|------|------| | `notification:create` | socialService, moderationService, galleryService, gateway(admin) | notificationService | 3次 | 创建通知(替换同步 RPC) | | `notification:push` | notificationService(self) | notificationService | 5次 | UniPush 推送(替换裸 goroutine) | ### 6.2 收益类 | Task Type | 生产者 | 消费者 | 重试 | 说明 | |-----------|--------|--------|------|------| | `revenue:exhibition` | galleryService | taskService | 3次 | 展示收益结算 | | `revenue:like-bet` | galleryService | taskService | 3次 | 点赞押注收益计算 | ### 6.3 定时任务类(Asynq Scheduler) | Task Type | 调度表达式 | 消费者 | 说明 | |-----------|-----------|--------|------| | `task:daily-reset` | `0 5 * * *` | taskService | 每日 05:00 重置任务 | | `asset:season-reset` | 按赛季配置 | assetService | 赛季重置 | | `statistic:materialize` | `*/5 * * * *` | statisticService | 物化视图刷新 | | `statistic:weekly-income` | `0 2 * * *` | statisticService | 周收入更新 | | `statistic:level-up` | `*/30 * * * *` | statisticService | 等级提升更新 | | `statistic:partition-create` | `5 0 * * *` | statisticService | events 表每日分区创建(替代 partitioner goroutine) | | `statistic:partition-drop` | `30 0 * * *` | statisticService | 过期分区清理(替代 partitioner goroutine) | ### 6.4 业务类 | Task Type | 生产者 | 消费者 | 重试 | 说明 | |-----------|--------|--------|------|------| | `moderation:auto-hide` | moderationService | moderationService | 3次 | 自动隐藏内容 | | `gallery:exhibition-expire` | galleryService | galleryService | 1次 | 展览自然到期(Asynq ProcessAt 精确时刻触发);ZSET + 每日兜底扫描作为补偿 | | `gallery:exhibition-settled` | galleryService | galleryService | 3次 | 统一结算任务(自然到期/手动下架/踢走等任意触发场景);handler 内做幂等 + 路由到收益计算/累计时长任务 | | `user:accumulate-hours` | galleryService | userService | 3次 | 增加用户累计上架时长(含手动下架,补回此前遗漏的;handler 内幂等检查 exhibition_id) | | `asset:accumulate-hours` | galleryService | assetService | 3次 | 增加资产累计展出时长(推动资产等级升级;幂等键 exhibition_id) | | `gallery:cleanup-display-status` | (scheduler) | galleryService | 0次 | display_status 不一致修复(每小时) | | `gallery:expired-exhibition-fallback` | (scheduler) | galleryService | 0次 | DB 兜底扫描过期展览(每天 04:00,补偿 Asynq delay task 遗漏) | | `starbook:collection-create` | starbookService | starbookService | 3次 | 收藏集创建异步处理 | | `aichat:chat` | aiChatService | aiChatService | 3次 | AI 对话异步处理 | | `cache:invalidate` | socialService | socialService | 1次 | 缓存失效 | --- ## 七、各服务改造方案 ### 7.1 notificationService — 核心消费者 🔴 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/consumer.go` | 注册 `notification:create` 和 `notification:push` handler(通过 `adapter.TaskConsumer.RegisterTask`) | | 新增 | `mq/producer.go` | 封装 `EnqueuePushPayload()`,只 import `adapter` | | 修改 | `main.go` | `mq.RegisterHandlers()` + `mq.StartConsumers()` | | 修改 | `service/notification_service.go` | 删除 `go func()` (L231),改为 `mq.EnqueuePushPayload(notif)` | **Handler 逻辑:** ``` notification:create (通过 adapter.TaskConsumer 注册) → 参数校验 → 事务内写 notifications + stats → mq.EnqueuePushPayload(notif) ← 走 adapter,不直接调 asynq notification:push (通过 adapter.TaskConsumer 注册) → 拉取用户活跃 cids → UniPush.Send() → 失败自动重试(最多5次) → 最终失败进死信队列 ``` > `CreateNotification` gRPC 接口保留,供 admin 面板通过 gateway 调用。 ### 7.2 galleryService — 核心生产者 🔴 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/consumer.go` | 注册 `gallery:exhibition-expire`(自然到期)+ `gallery:exhibition-settled`(通用结算)handler(通过 `adapter.TaskConsumer.RegisterTask`) | | 新增 | `mq/producer.go` | 封装事件发布 + task 入队 + `EnqueueExhibitSettled`,只 import `pkg/mq/adapter`,不直接接触 asynq/streams | | 修改 | `service/exhibition_service.go` | `RemoveFromSlot()` 和 `RemoveExhibitionByAsset()` 末尾改为 `mq.EnqueueExhibitSettled(exhibition_id, source)`,不再直接调 RPC | | 修改 | `main.go` | `mq.RegisterHandlers()` + `mq.StartConsumers()` | | **删除** | `service/cleanup_worker.go` | 整个文件:轮询 → Asynq delay task + scheduler | | 新增 | Asynq scheduler | 注册 `gallery:cleanup-display-status`(每小时),替代 `cleanupInvalidDisplayStatus()` | | 新增 | Asynq scheduler | 注册 `gallery:expired-exhibition-fallback`(每天 04:00),DB 兜底查询过期展览,补偿 Asynq delay task 可能遗漏的 | **展览生命周期(统一结算模型,所有 MQ 调用走 adapter 层):** ``` [展览上架 — 通过 adapter] → 写 ZSET(保留兼容) → mq.EnqueueExhibitionExpire(exhibition_id, expireAt) └→ adapter.Get().TaskProducer().EnqueueAt( adapter.Task{Type: "gallery:exhibition-expire", MaxRetry: 3, ...}, expireAt) [展品结算触发 — 任何路径统一入口] ★ 核心改进 所有下架 / 清扫代码路径都调: mq.EnqueueExhibitSettled(exhibition_id, source) source ∈ {"natural","manual","kick"} └→ adapter.Get().TaskProducer().Enqueue( adapter.Task{Type: "gallery:exhibition-settled", MaxRetry: 3, Payload: ...}) [gallery:exhibition-settled Handler — 三场景统一处理] → 幂等检查: 查 exhibition.settled → 已结算则 return → 标记 settled=true(不再仅看 processed,覆盖全场景) → mq.PublishExhibitionSettled(e) → stream:exhibition XADD → EnqueueRevenueExhibition(payload) → EnqueueRevenueLikeBet(payload) → EnqueueAccumulateUserHours(exhibition_id, hours) ← ★ 补回手动下架时长 → EnqueueAccumulateAssetHours(exhibition_id, hours) [gallery:exhibition-expire Handler — 自然到期精确触发] → 调 assetClient.GetAssetLikeCount(assetID) RPC 查点赞数 → 调 repo 计算 actual_hours = (expireAt - startTime) / 3600000 → mq.EnqueueExhibitSettled(exhibition_id, source="natural") > 幂等:entrance 查 settled 标志 + revenue/累计时长 handler 内各自去重(exhibition_id 幂等键) ``` ### 7.3 taskService — 收益计算 + 定时任务 🔴 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/consumer.go` | 注册 `revenue:exhibition`、`revenue:like-bet` handler + `task:daily-reset` scheduler(通过 `adapter` 接口注册) | | 新增 | `mq/producer.go` | 封装 Streams 发布(统计埋点),只 import `adapter` | | 修改 | `main.go` | `mq.RegisterHandlers()` + 启动 Asynq Server + Streams Consumer | | **删除** | `worker/daily_reset_worker.go` | 整个文件 | | 修改 | `service/revenue_service.go` | `OnExhibitionCompleted` 和 `RecordLikeBetRevenue` 改为 handler 内部函数,逻辑不变 | **Handler 逻辑:** ``` revenue:exhibition (payload: {exhibition_id, asset_id, slot_id, occupier_uid, occupier_star_id, slot_owner_uid, start_time, expire_at, like_count}) → 幂等检查: SELECT 是否已有同 exhibition_id 的 revenue record → 已存在则 return → 调 assetLevelService.CalculateRevenue(assetID, likeCount, startTime, expireAt, 0) → CreateRevenueRecord → 写 DB → 调 assetLevelService.AddExhibitionHours → 失败仅日志,不重试全任务 → 调 userRPCClient.AddExhibitionHours → 失败仅日志,不重试全任务 → 失败重试3次(幂等检查保证重复调用安全) revenue:like-bet (payload: {exhibition_id, asset_id, start_time, expire_at}) → 查 exhibition 下所有 asset_likes → 按 bet_order 计算每笔金额 → 批量写 like_bet_revenue_records → DB 唯一约束 uk_like_bet_unique(exhibition_id, like_id) 保证幂等 → 失败重试3次 → 死信 ``` ### 7.4 socialService — 替换同步 RPC 🟡 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/producer.go` | 封装 `notification:create` 入队 + `stream:social` 发布 | | 修改 | `service/asset_like_service.go` | L142 `fireLikeNotification` 改为 `mq.EnqueueLikeNotification(...)`(走 adapter) | | 可选删除 | `client/notification_client.go` | 不再需要同步 RPC 调 notificationService | **点赞流程改造后:** ``` LikeAsset() → assetClient.LikeAsset() RPC ✅ 保持同步(返回 likeCount) → XADD stream:social {type: "asset.like", ...} ← 替代 statistic.TrackEvent() → mq.EnqueueLikeNotification(...) ← 替代同步 RPC 调 notificationService(走 adapter) ``` > `statistic.TrackEvent()` 调用改为 XADD stream:social,由 statisticService 的 Streams consumer 消费。channel_sink 被删除后 TrackEvent() 不再可用。 ### 7.5 moderationService — 替换同步 RPC 🟡 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/producer.go` | 封装 `notification:create` 入队 | | 新增 | `mq/consumer.go` | 注册 `moderation:auto-hide` handler | | 修改 | `client/notification_client.go` | 改为 `mq.EnqueueReportNotice(...)`(删除 `go func()` fire-and-forget goroutine) | | 修改 | `service/report_service.go` | L185 `go func() { s.notifClient.SendXxx() }()` 改为 `mq.EnqueueReportNotice(...)` | ### 7.6 statisticService — channel → Streams 🟡 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/consumer.go` | Streams 消费者组 + Asynq scheduler | | 修改 | `main.go` | 启动 Streams Consumer + Asynq Server | | **删除** | `sink/event_sink.go` | EventSink 接口,随 channel_sink 一起废弃 | | **删除** | `sink/channel_sink.go` | 内存 channel → 生产者直接 XADD Stream | | **保留迁移** | `worker/partitioner.go` | → Asynq scheduler(`statistic:partition-create` 00:05 + `statistic:partition-drop` 00:30);events 表每日分区管理不能停 | | **删除** | `worker/event_flusher.go` | goroutine 批量写入 → Streams consumer | | **删除** | `worker/materializer.go` | ticker → Asynq scheduler | | **删除** | `worker/metric_weekly_user_income_updater.go` | ticker → Asynq scheduler | | **删除** | `worker/metric_upcoming_level_ups_updater.go` | ticker → Asynq scheduler | **改造后架构:** ``` [消费] Redis Streams Consumer Group stream:user / stream:asset / stream:social / stream:exhibition / stream:moderation → 批量写入 event_partitions → 触发增量更新 [消费] Asynq Scheduler statistic:materialize (每5分钟) → REFRESH MATERIALIZED VIEW statistic:weekly-income (每天) → 周收入更新 statistic:level-up (每30分钟) → 等级提升更新 ``` ### 7.7 userService — 事件发布 🟢 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/producer.go` | 封装 `stream:user` 发布 | | 新增 | `mq/consumer.go` | 注册 `user:accumulate-hours` handler(接收 galleryService 派发的累计时长任务,含手动下架补漏) | | 修改 | service 层 | Register/UpdateProfile 末尾追加 XADD | ### 7.8 assetService — 定时任务 + 事件 🟢 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/producer.go` | 封装 `stream:asset` 发布 | | 新增 | `mq/consumer.go` | 注册 `asset:season-reset` scheduler + `asset:accumulate-hours` handler(接收 galleryService 派发的资产累计时长) | | **删除** | `worker/season_reset_worker.go` | cron → scheduler | | 修改 | `service/mint_service.go` | Mint 完成后 XADD stream:asset | ### 7.9 aiChatService — 异步对话 🟢 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/producer.go` | Chat 时 Enqueue aichat:chat | | 新增 | `mq/consumer.go` | Handler 调 LLM → 写 DB | ### 7.10 starbookService — 异步收藏集 🟢 **改造内容:** | 改动 | 文件 | 说明 | |------|------|------| | 新增 | `mq/producer.go` | CreateCollection 时 Enqueue | | 新增 | `mq/consumer.go` | Handler 处理异步逻辑 | ### 7.11 activityService — 不改 ⚪ 已有 Redis Streams combo worker,只把 key 统一命名为 `stream:activity`。 ### 7.12 gateway — 微调 ⚪ - 新增 `stream:activity` 消费者组连接(已有 WebSocket hub,整合) - admin 发系统通知从同步 RPC 改为 `mq.EnqueueAdminNotification(...)`(走 adapter) --- ## 八、数据流对比 ### 改造前(展示收益 + 点赞押注) ``` galleryService.CleanupWorker (每分钟 ticker 轮询) → RPC assetService.GetAssetLikeCount() → 本地计算收益(硬编码 R0=5) → RPC taskService.OnExhibitionCompleted() ← 任务服务挂了收益就丢了 → RPC taskService.RecordLikeBetRevenue() ← 同上 → 失败只打 warn 日志,无重试 ``` ### 改造后 ``` [展览上架时 — 业务侧 mq.Xxx,broker 无关] mq.EnqueueExhibitionExpire(exhibition_id, expireAt) └→ adapter.Get().TaskProducer().EnqueueAt(...) [到期触发 — handler 也走 adapter.TaskConsumer] gallery:exhibition-expire handler → mq.PublishExhibitionExpired(e) → adapter.Get().EventProducer().Publish("stream:exhibition", e) → mq.EnqueueRevenueExhibition(payload) ← 重试 + 死信(broker 无关) → mq.EnqueueRevenueLikeBet(payload) ← 重试 + 死信 → 标记 exhibition.settled=true [taskService 消费 — 也通过 adapter.TaskConsumer 路由] revenue:exhibition handler → 计算收益 → 写DB revenue:like-bet handler → 批量写入 ``` **关键改进:** taskService 挂了不影响事件产生,恢复后自动从上次 ACK 位置继续处理。 --- ## 九、错误处理 ### 9.1 Asynq 重试策略 | 场景 | 重试次数 | 退避策略 | 最终失败 | |------|----------|----------|----------| | 通知创建 | 3 | 指数退避 10s/30s/90s | 死信队列 → 告警 | | 推送 | 5 | 指数退避 10s/30s/90s/270s/810s | 死信队列 → 告警 | | 收益计算 | 3 | 指数退避 10s/30s/90s | 死信队列 → 告警 | | 缓存失效 | 1 | 无 | 丢弃(缓存有 TTL 兜底) | | 定时任务 | 0 | N/A | 下次调度自动重试 | ### 9.2 Redis Streams 错误处理 - 消费者组 ACK 机制:消息处理后显式 ACK,未 ACK 的消息在 PEL 中 - 消费者崩溃重启:从上次 ACK 位置继续,消息不丢 - Stream 有最大长度限制(100000),防止内存无限增长 ### 9.3 死信处理 ``` [Asynq Dead Letter Queue] → 定时巡检(每小时) → 输出 ERROR 日志(含 task_type + payload + 错误原因) → 关键 task(收益类)触发钉钉/飞书告警 → 人工介入:通过 admin 接口重放或跳过 ``` --- ## 十、配置 ### 10.1 .env 新增配置 ```bash # Asynq ASYNQ_REDIS_ADDR=localhost:6379 ASYNQ_REDIS_DB=2 ASYNQ_REDIS_PASSWORD= ASYNQ_CONCURRENCY=10 # 并发处理数 # Redis Streams STREAMS_REDIS_ADDR=localhost:6379 STREAMS_REDIS_DB=3 STREAMS_CONSUMER_GROUP=topfans-service STREAMS_MAX_LEN=100000 ``` ### 10.2 go.mod 新增依赖 ``` github.com/hibiken/asynq # Asynq 任务队列 ``` Redis Streams 使用已有的 `github.com/redis/go-redis/v9`,不需要新包。 --- ## 十一、迁移计划 ### 阶段一:基础设施(1-2天) 1. 实现 `pkg/mq/asynq/` — client、server、tasks 常量、middleware 2. 实现 `pkg/mq/streams/` — producer、consumer、keys 常量 3. 实现 `pkg/mq/config.go` 4. 单元测试 ### 阶段二:核心服务改造(3-5天) 1. **notificationService** — Asynq handler(通知创建 + 推送) 2. **galleryService** — 替换 CleanupWorker 3. **taskService** — 替换 DailyResetWorker + 接收 revenue task ### 阶段三:次级服务改造(2-3天) 4. **socialService** — 异步通知 5. **moderationService** — 异步通知 + auto-hide 6. **statisticService** — channel → Streams + ticker → scheduler ### 阶段四:新增事件能力(1-2天) 7. **userService** — 事件发布 8. **assetService** — 事件发布 + scheduler 9. **aiChatService** — 异步对话 10. **starbookService** — 异步收藏集 ### 阶段五:收尾(1天) 11. **gateway** — Streams consumer + admin 异步通知 12. **activityService** — 统一 key 命名 13. 全链路压力测试 14. 清理删除的旧代码 --- ## 十二、函数/接口迁移结构 ### 12.1 收益计算相关函数去向 | 函数 | 当前位置 | 改造后 | 说明 | |------|----------|--------|------| | `cleanup_worker.calculateExhibitionRevenue` | galleryService | **删除** | R0 硬编码 5 作为兜底值传给 taskService;改造后 handler 直接调 `assetLevelService.CalculateRevenue()`(R0 按等级从 DB 读),无需兜底 | | `cleanup_worker.calculateExhibitionRevenue` | galleryService | **删除** | R0 硬编码 5 作为兜底值传给 taskService;改造后 handler 直接调 `assetLevelService.CalculateRevenue()`(R0 按等级从 DB 读),无需兜底 | | `cleanup_worker.cleanup` | galleryService | **删除** | ticker 轮询 → Asynq delay task | | `cleanup_worker.cleanupExpiredExhibitions` | galleryService | **删除** | ZSET+DB 轮询 → Asynq delay task | | `cleanup_worker.cleanupAssetsFromZSET` | galleryService | **删除** | ZSET 按 asset 逐个清理 → Asynq delay task | | `cleanup_worker.cleanupExpiredExhibitionsFromDB` | galleryService | **删除** | DB 兜底轮询 → Asynq delay task | | `cleanup_worker.cleanupInvalidDisplayStatus` | galleryService | → Asynq scheduler `gallery:cleanup-display-status` | display_status 不一致修复,独立定时任务(每小时) | | `revenueService.OnExhibitionCompleted` | taskService | → `mq/consumer.go` handler 内部函数 | 改为 `revenue:exhibition` handler,逻辑复用 | | `revenueService.RecordLikeBetRevenue` | taskService | → `mq/consumer.go` handler 内部函数 | 改为 `revenue:like-bet` handler,逻辑复用 | | `revenueService.CalculateExhibitionRevenue` | taskService | **保留** | 参考实现(未被调用),R0 硬编码 5;保留供测试对比 | | `revenueService.CalculateBuff` | taskService | **保留** | 纯函数,taskService 包内使用 | | `CalculateLikeBetRevenue` | taskService | **保留不动** | 纯计算函数,handler 内部调用 | | `assetLevelService.CalculateRevenue` | assetService | **保留不动** | 唯一正式版本,R0 从 `levelConfig.HourlyRevenue`(DB)读取 | | `assetLevelService.CalculateBuff` | assetService | **保留不动** | 纯函数,assetService 包内使用(与 taskService.CalculateBuff 逻辑相同但各自独立) | ### 12.2 AssetLevelService 接口 | 项目 | 改造前 | 改造后 | |------|--------|--------| | 接口定义 | taskService `revenue_service.go` | **保持不变**,taskService 仍需调用 | | 实现方 | assetService `asset_level_service.go` | **保持不变** | | 调用方式 | galleryService RPC → taskService → assetService RPC | Asynq handler → assetService RPC | | `CalculateRevenue` | OnExhibitionCompleted 内调用 | revenue:exhibition handler 内调用,逻辑不变 | | `AddExhibitionHours` | OnExhibitionCompleted 内调用 | revenue:exhibition handler 内调用,逻辑不变 | ### 12.3 定时任务 Worker 迁移 | Worker | 当前位置 | 改造后 | |--------|----------|--------| | `CleanupWorker` | galleryService | **删除文件** → `gallery:exhibition-expire` delay task + `gallery:cleanup-display-status` scheduler | | `DailyResetWorker` | taskService | **删除文件** → `task:daily-reset` Asynq scheduler | | `SeasonResetWorker` | assetService | **删除文件** → `asset:season-reset` Asynq scheduler | | `Materializer` | statisticService | **删除文件** → `statistic:materialize` Asynq scheduler | | `MetricWeeklyUserIncomeUpdater` | statisticService | **删除文件** → `statistic:weekly-income` Asynq scheduler | | `MetricUpcomingLevelUpsUpdater` | statisticService | **删除文件** → `statistic:level-up` Asynq scheduler | | `EventFlusher` | statisticService | **删除文件** → Redis Streams consumer | | `Partitioner` | statisticService | **保留文件,触发方式迁移** → `statistic:partition-create` + `statistic:partition-drop` Asynq scheduler | ### 12.4 缓存/状态层迁移 | 组件 | 当前位置 | 改造后 | |------|----------|--------| | `sink.EventSink` 接口 | statisticService | **删除文件** → 统一用 `streams.Producer` | | `sink.ChannelEventSink` | statisticService | **删除文件** → 生产者直接 XADD Stream | | `worker.Partitioner` | statisticService | **保留文件** → Asynq scheduler 触发(events 表每日分区管理不能停) | | `database.GetExpiredAssets` (ZSET) | pkg/database | **保留** — 兼容,Asynq delay task 是主路径 | | `database.RemoveExpiringAsset` | pkg/database | **保留** — Asynq handler 内调用 | --- ## 十三、风险与缓解 | 风险 | 严重度 | 缓解 | |------|--------|------| | Asynq 依赖 Redis,Redis 挂了全挂 | 高 | Redis 已有主从,增加哨兵保活;关键同步路径保留 RPC 兜底 | | `gallery:exhibition-expire` handler 失败导致收益漏算 | **高** | 重试 1 次 + 保留 ZSET 每日兜底扫描(原 `cleanupExpiredExhibitionsFromDB` 逻辑改为 Asynq scheduler 每日跑一次) | | `revenue:exhibition` 重试导致重复创建收益记录 | **高** | handler 入口先按 `exhibition_id` 查已有记录,已存在直接返回 + 后续可加 DB 唯一约束 | | Streams 消费者组 rebalance 导致短暂不可用 | 中 | 消费者组启动时有 5s 的 claim pending 消息逻辑 | | 异步化导致前端感知延迟 | 低 | 通知/收益创建本身就有 DB 写入延迟,异步化不增加用户可见延迟 | | 收益计算逻辑迁移出错 | 中 | `revenue:exhibition` handler 内部完全复用现有 `OnExhibitionCompleted` 逻辑,只改调用方式 |