# Worker 队列迁移变更文档:DB 轮询 → Redis Stream 分发 > 状态:草案 > 日期:2026-07-10 > 范围:`internal/worker/`、`internal/app/`、`internal/config/`、部署配置 > 原则:Redis 负责消息分发(XADD/XREADGROUP),DB 只负责 job 元数据落地与审计追踪 --- ## 1. 背景与动机 ### 1.1 当前架构 Worker 子系统由 `internal/worker/worker.go` 的 `WorkerPool` 实现,核心行为: - **入队**:`Enqueue()` 将 `BackgroundJob` 记录 INSERT 到 PostgreSQL `background_jobs` 表 - **消费**:每个 worker goroutine 每 500ms 轮询 DB,`SELECT ... FOR UPDATE SKIP LOCKED` claim 最早到期的 job - **执行**:`perform()` 调用注册的 `JobHandler`,完成后 UPDATE 状态为 `completed`/`retrying`/`dead` - **补偿**:`RequeueStaleJobs()` 扫描超过 `staleLockTimeout`(默认 15 分钟)的 `running` job,重置为 `retrying` Redis 在项目中已用于三种场景,但 **均不涉及 job 队列**: | 角色 | 实现 | 代码位置 | |------|------|----------| | WS 跨实例广播 | Redis Pub/Sub `PSubscribe` | `internal/ws/broadcast.go` | | 事件总线 | Watermill Redis Streams (consumer group) | `internal/pubsub/redis_pubsub.go` | | 通知投递 | Watermill Router + Redis Streams | `internal/service/notification_delivery_service.go` | ### 1.2 当前架构的问题 1. **轮询开销**:500ms 固定轮询,空闲时持续消耗 DB 连接和 CPU 2. **延迟下限**:job 从入队到执行至少 0–500ms 延迟 3. **DB 连接占用**:每个 worker goroutine 持续占用一个 DB 连接做轮询 4. **concurrency 配置未生效**:`Bootstrap()` 使用 `NewWorkerPool(db)` 而非 `NewWorkerPoolWithOptions(db, WithWorkerCount(n))`,实际只有 1 个 goroutine 5. **基础设施不统一**:Redis 已有 watermill-redisstream + go-redis,但 job 队列绕过 Redis 走 DB 轮询 ### 1.3 迁移目标 ``` Enqueue(): 1. INSERT into background_jobs (status=queued) ← DB 落地,拿到 job ID 2. XADD job stream {job_id, job_type, queue} ← Redis 分发 Worker: 1. XREADGROUP BLOCK → 拿到 job_id ← Redis 阻塞拉取(零轮询) 2. SELECT * FROM background_jobs WHERE id=? ← DB 加载 payload + 元数据 3. UPDATE status=running, locked_by=... ← DB 标记执行中 4. handler(ctx, job) ← 执行 5. UPDATE status=completed/retrying/dead ← DB 更新最终状态 6. XACK ← Redis 确认 Sweep (补偿): 每 30s 扫描 DB 中 status=queued/retrying 且 scheduled_at <= now 的 job 重新 XADD,兜底 Redis 投递失败的情况 ``` --- ## 2. 影响面分析 ### 2.1 Job Type 全量清单(37 个) 迁移不改变任何 job type 定义,仅改变 `Enqueue` 和 `run` 的内部实现。 | Job Type | Queue | 调用方 | 是否延迟 | 幂等键 | |----------|-------|--------|----------|--------| | `message:send_reply` | `high` | `service.EnqueueSendReply` | 否 | `message:send_reply:{msgID}` | | `event:dispatch_async` | `events` | `channel.Dispatcher.DispatchAsync` | 否 | 无 | | `event:listener_dispatch` | `events` | `dispatch.EventDispatcher` | 否 | 无 | | `automation:team_email_delivery` | `automation` | `automation.ActionService` | 否 | 无 | | `automation:transcript_delivery` | `automation` | `automation.ActionService` | 否 | 无 | | `automation:webhook_delivery` | `automation` | `automation.ActionService` | 否 | 无 | | `automation:macro_execution` | `medium` | `automation.MacroService` | 否 | 无 | | `csat:survey_send` | `automation` | `automation.CsatSurveyListener` | 否 | `csat-survey:{convID}` | | `csat:template_create` | `automation` | `service.CsatTemplateService` | 否 | `csat-template:{tplID}:{ts}` | | `scheduled:trigger_items` | `scheduled_jobs` | `service.EnqueueScheduledItemsTrigger` | **是**(1h 周期) | `scheduled:trigger_items:{bucket}` | | `campaign:trigger_oneoff` | `low` | `conversationMaintenanceRunner` | 否 | `campaign:trigger_oneoff:{id}` | | `conversation:reopen_snoozed` | `low` | `conversationMaintenanceRunner` | 否 | 无 | | `account:conversations_resolution_scheduler` | `scheduled_jobs` | `conversationMaintenanceRunner` | 否 | 无 | | `conversation:resolution` | `low` | `conversationMaintenanceRunner` | 否 | 无 | | `conversation:update_message_status` | `deferred` | `service.EnqueueConversationMessageStatusUpdate` | **是**(可选) | `conversation:update_message_status:{convID}:{status}:{ts}` | | `conversation:bulk_action` | `medium` | `service.EnqueueConversationBulkAction` | 否 | 无 | | `contact:bulk_action` | `medium` | `service.EnqueueContactBulkAction` | 否 | 无 | | `conversation:delete_object` | `low` | `service.ConversationService` | 否 | `conversation-delete:{acctID}:{convID}` | | `contact:import` | `low` | `service.ContactService` | 否 | `contact-import:{id}` | | `contact:export` | `low` | `service.ContactService` | 否 | `contact-export:{id}` | | `search:index` | `search` | `service.DurableSearchIndexer` | 否 | 无 | | `sla:trigger_accounts` | `scheduled_jobs` | `service.EnqueueSlaAccountsScan` | **是**(5min 周期) | `sla:trigger_accounts:{bucket}` | | `sla:process_account` | `medium` | `slaProcessingJobRunner` | 否 | 无 | | `sla:process_applied` | `medium` | `slaProcessingJobRunner` | 否 | 无 | | `captain:copilot_response` | `default` | `service.CopilotService` | 否 | `captain:copilot_response:{msgID}` | | `captain:conversation_response_builder` | `default` | `service.CaptainConversationService` | **是**(可选,附件等待) | `captain:conversation_response_builder:message:{msgID}` | | `captain:document_sync` | `low` | `service.CaptainDocumentService` | 否 | `captain:document_sync:auto:...`(自动同步时) | | `captain:document_crawl` | `low` | `service.CaptainDocumentService` | 否 | `captain:document_crawl:{acctID}:{id}` | | `captain:document_page_crawl_parse` | `low` | `service.CaptainDocumentService` | 否 | `captain:document_page_crawl_parse:...` | | `captain:document_response_builder` | `low` | `service.CaptainDocumentService` | 否 | 无 | | `captain:documents_schedule_syncs` | `scheduled_jobs` | `service.EnqueueCaptainDocumentScheduleSyncs` | **是**(24h 周期) | `captain:documents_schedule_syncs:{bucket}` | | `captain:llm_update_embedding` | `low` | `service.CaptainDocumentService` | 否 | 无 | | `captain:article_translate` | `low` | `service.ArticleService` | 否 | `captain:article_translate:...` | | `inbox:sync_templates` | `low` | `service.InboxService` | 否 | 无 | | `webhook:incoming_message_persist` | `default`/`low` | `webhook.IncomingPersister` | 否 | `webhook:incoming_message:{inboxID}:{sourceID}` | | `webhook:message_status_update` | `low` | `webhook.IncomingPersister` | 否 | `webhook:message_status:...` | | `webhook:contact_messages_status_update` | `low` | `webhook.IncomingPersister` | 否 | `webhook:contact_messages_status:...` | | `webhook:webwidget_triggered` | `default` | `service.WidgetService` | 否 | 无 | | `reporting:rollup_day` | `low` | `service.EnqueueReportingRollupDay` | 否 | `reporting:rollup_day:{acctID}:{date}` | ### 2.2 Queue 名称汇总(8 个) | Queue | 用途 | 延迟 job | 优先级特征 | |-------|------|----------|------------| | `default` | 通用 | 否 | 混杂 | | `high` | 消息发送回复 | 否 | 最高优先 | | `medium` | 批量操作、SLA 处理 | 否 | 中等 | | `low` | 文档同步、爬取、导入导出、删除 | 否 | 可延迟 | | `events` | 事件分发 | 否 | 需低延迟 | | `automation` | 自动化规则执行 | 否 | 需及时 | | `search` | 搜索索引 | 否 | 可延迟 | | `scheduled_jobs` | 定时周期任务 | **是** | 按时间触发 | | `deferred` | 延迟任务 | **是** | 按时间触发 | | `purgable` | 可清理 | 否 | 低优先级 | 迁移后每个 queue 对应一个 Redis Stream key:`gochat:jobs:{queue}`。 ### 2.3 延迟 Job(scheduled_at)的使用场景 5 个 job type 使用 `WithScheduledAt`,需要特殊处理: 1. **`scheduled:trigger_items`**:1 小时周期,自我重新入队 2. **`sla:trigger_accounts`**:5 分钟周期,自我重新入队 3. **`captain:documents_schedule_syncs`**:24 小时周期,自我重新入队 4. **`conversation:update_message_status`**:可选延迟 5. **`captain:conversation_response_builder`**:可选延迟(等待附件上传完成,1–5s) 延迟 job 在 Redis Stream 中不能直接投递(Stream 无原生延迟语义),需要延迟投递机制(见 §3.4)。 ### 2.4 Priority 使用情况 `WithPriority` 已实现但**生产代码中无任何调用方**。仅 `worker_test.go` 测试中使用了 `WithPriority(10)`。 迁移后 priority 语义改为 queue 级别优先级(worker 先读 high queue 再读 low),不再需要 DB `ORDER BY priority DESC`。`BackgroundJob.Priority` 字段保留供 DB 审计,但不再影响消费顺序。 ### 2.5 Idempotency Key 使用情况 19 个 job type 使用了 `WithIdempotencyKey`,依赖 DB 唯一索引 `idx_background_jobs_idempotency_key_unique` 做去重。迁移后 idempotency 仍在 DB 层保证,Redis 不参与去重。 --- ## 3. 详细设计 ### 3.1 WorkerPool 结构体变更 **文件**:`internal/worker/worker.go` ```go type WorkerPool struct { db *gorm.DB rdb redis.UniversalClient // 新增 handlers map[string]JobHandler queues []string workerID string workerCount int pollInterval time.Duration // 保留,仅 fallback 模式使用 staleLockTimeout time.Duration backoff BackoffFunc now func() time.Time // 新增:Redis 队列参数 streamPrefix string // 默认 "gochat:jobs" consumerGroup string // 默认 "gochat-workers" blockTimeout time.Duration // XREADGROUP block 时长,默认 5s sweepInterval time.Duration // 补偿扫描间隔,默认 30s mu sync.RWMutex ctx context.Context cancel context.CancelFunc wg sync.WaitGroup } ``` 新增 Option 函数: ```go func WithRedisClient(rdb redis.UniversalClient) Option func WithStreamPrefix(prefix string) Option func WithConsumerGroup(group string) Option func WithBlockTimeout(d time.Duration) Option func WithSweepInterval(d time.Duration) Option ``` ### 3.2 Enqueue 变更 **当前流程**: 1. 构建 `BackgroundJob` 结构 2. 处理 idempotency key 去重(DB 查询) 3. `db.Create(job)` 插入 DB 4. 返回 job **迁移后流程**: 1. 构建 `BackgroundJob` 结构(不变) 2. 处理 idempotency key 去重(DB 查询,不变) 3. `db.Create(job)` 插入 DB(不变) 4. **新增**:如果 `scheduled_at <= now`,`XADD` 到 Redis Stream 5. **新增**:如果 `scheduled_at > now`,跳过 XADD(由 sweep 机制在到期时投递) 6. 返回 job(不变) ```go func (wp *WorkerPool) Enqueue(ctx context.Context, jobType string, payload any, opts ...EnqueueOption) (*model.BackgroundJob, error) { // ... 现有逻辑:构建 job、idempotency、DB 插入(完全不变)... // 新增:DB 插入成功后投递到 Redis if wp.rdb != nil && !job.ScheduledAt.After(wp.now()) { if err := wp.pushToRedis(ctx, job); err != nil { // Redis 投递失败不影响 DB 已提交的 job // sweep 补偿机制会重新投递 applogger.L().Warnf( "redis XAdd failed for job %d (queue=%s), will be picked up by sweep: %v", job.ID, job.Queue, err, ) } } // scheduled_at > now 的 job 不立即投递 // sweep goroutine 会定期扫描到期 job 并投递 return job, nil } func (wp *WorkerPool) pushToRedis(ctx context.Context, job *model.BackgroundJob) error { streamKey := wp.streamKeyFor(job.Queue) return wp.rdb.XAdd(ctx, &redis.XAddArgs{ Stream: streamKey, Values: map[string]interface{}{ "job_id": fmt.Sprintf("%d", job.ID), "job_type": job.JobType, "queue": job.Queue, }, }).Err() } func (wp *WorkerPool) streamKeyFor(queue string) string { return fmt.Sprintf("%s:%s", wp.streamPrefix, queue) } ``` **影响分析**: - 所有 `Enqueue*` 辅助函数(`EnqueueSendReply`、`EnqueueScheduledItemsTrigger` 等)**不需要改动**,它们调用 `wp.Enqueue()` 的签名和返回值不变 - `SetWorkerPool` 调用方**不需要改动** - idempotency 逻辑完全不变,仍在 DB 层 ### 3.3 Worker 消费循环变更 **当前**:`run()` 每 500ms 调用 `ProcessOne()` → `claimNext()` → DB `SELECT ... FOR UPDATE SKIP LOCKED` **迁移后**:`run()` 使用 `XREADGROUP BLOCK` 阻塞拉取 ```go func (wp *WorkerPool) run(ctx context.Context) { defer wp.wg.Done() if wp.rdb == nil { // fallback:Redis 不可用时回退到 DB 轮询模式 wp.runDBPollLoop(ctx) return } streams := wp.streamKeys() // 监听所有已配置 queue 的 stream for { if ctx.Err() != nil { return } // XREADGROUP BLOCK results, err := wp.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{ Group: wp.consumerGroup, Consumer: wp.workerID, Streams: streams, Count: 1, Block: wp.blockTimeout, // 5s }).Result() if ctx.Err() != nil { return } if err != nil && !errors.Is(err, redis.Nil) { applogger.L().Errorf("XReadGroup error: %v", err) time.Sleep(time.Second) // 退避 continue } for _, xstream := range results { for _, msg := range xstream.Messages { wp.processRedisMessage(ctx, msg) } } } } func (wp *WorkerPool) processRedisMessage(ctx context.Context, msg redis.XMessage) { jobIDStr, ok := msg.Values["job_id"] if !ok { // 格式错误,直接 ACK 丢弃 wp.ackRedis(ctx, msg.Stream, msg.ID) return } jobID, err := strconv.ParseUint(jobIDStr.(string), 10, 64) if err != nil { wp.ackRedis(ctx, msg.Stream, msg.ID) return } // 从 DB 加载 job var job model.BackgroundJob if err := wp.db.WithContext(ctx).First(&job, jobID).Error; err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { // job 已被删除,ACK 丢弃 wp.ackRedis(ctx, msg.Stream, msg.ID) return } // DB 错误,不 ACK,让 Redis 重投给其他 consumer applogger.L().Errorf("load job %d from DB failed: %v", jobID, err) return } now := wp.now() if err := archiveExhaustedRetries(wp.db.WithContext(ctx).Where("id = ?", jobID), now); err != nil { applogger.L().Errorf("archive exhausted job %d failed: %v", jobID, err) return } // Claim 检查:防止多 consumer 竞争同一 job result := wp.db.WithContext(ctx).Model(&model.BackgroundJob{}). Where("id = ? AND status IN ? AND scheduled_at <= ? AND attempts < max_attempts", jobID, []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, now, ). Updates(map[string]any{ "status": model.BackgroundJobStatusRunning, "locked_at": &now, "locked_by": wp.workerID, "attempts": gorm.Expr("attempts + 1"), }) if result.Error != nil { applogger.L().Errorf("claim job %d failed: %v", jobID, result.Error) return } if result.RowsAffected == 0 { // 已被其他 consumer claim 或尚未到期,ACK 跳过 wp.ackRedis(ctx, msg.Stream, msg.ID) return } // 重新加载完整 job(含 attempts 等更新后的字段) wp.db.WithContext(ctx).First(&job, jobID) // 执行 handler if err := wp.perform(ctx, &job); err != nil { wp.fail(ctx, &job, err) } // 确认 Redis 消息 wp.ackRedis(ctx, msg.Stream, msg.ID) } func (wp *WorkerPool) ackRedis(ctx context.Context, stream, msgID string) { if err := wp.rdb.XAck(ctx, stream, wp.consumerGroup, msgID).Err(); err != nil { applogger.L().Warnf("XAck failed for stream=%s msgID=%s: %v", stream, msgID, err) } } ``` **影响分析**: - `perform()` 和 `fail()` 方法**不变** - `ProcessOne()` 方法保留,仅用于测试和 fallback 模式 - `claimNext()` 保留,仅用于 DB 轮询 fallback 路径 - `JobHandler` 签名 `func(ctx, *model.BackgroundJob) error` **不变** ### 3.4 延迟 Job 投递机制 **问题**:Redis Stream 不支持延迟投递。`scheduled_at > now` 的 job 入队时不应立即 XADD。 **方案**:Sweep goroutine 定期扫描到期 job,投递到 Redis。 ```go func (wp *WorkerPool) sweepLoop(ctx context.Context) { defer wp.wg.Done() ticker := time.NewTicker(wp.sweepInterval) // 30s defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if err := wp.sweepDueJobs(ctx); err != nil { applogger.L().Warnf("sweep due jobs failed: %v", err) } } } } func (wp *WorkerPool) sweepDueJobs(ctx context.Context) error { now := wp.now() var jobs []model.BackgroundJob // 查找已到期但仍在 queued/retrying 状态的 job if err := wp.db.WithContext(ctx). Where("status IN ? AND scheduled_at <= ?", []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, now, ). Limit(100). Find(&jobs).Error; err != nil { return err } for _, job := range jobs { // 幂等投递:重复 XADD 不会导致重复执行 // 因为 processRedisMessage 中有 DB claim 检查 if err := wp.pushToRedis(ctx, &job); err != nil { applogger.L().Warnf("sweep: XAdd failed for job %d: %v", job.ID, err) } } return nil } ``` **幂等保证**:sweep 可能对同一个 job 多次 XADD(因为不知道 Redis 中是否已有该 job 的待处理消息)。这不会导致重复执行,因为 `processRedisMessage` 中 DB claim 检查保证只有一个 consumer 能成功将 status 从 `queued`/`retrying` 改为 `running`。 **影响分析**: - 延迟 job 的 `Enqueue` 调用方**不需要改动** - 延迟从最坏 500ms(轮询间隔)变为最坏 30s(sweep 间隔) - 对于需要更精确延迟的 job(如 `captain:conversation_response_builder` 的 1–5s 附件等待),可配置更短的 sweep 间隔或单独处理 ### 3.5 RequeueStaleJobs 变更 **当前**:扫描 `status=running AND locked_at < cutoff`,重置为 `retrying`。 **迁移后**:仍由 `Start()` 在启动时调用一次(回收上次崩溃时正在执行的 job),并在 sweep 循环中定期调用。stale job 未耗尽 `max_attempts` 时转为 `retrying`;已耗尽时原子转为 `dead` 并写入 `failed_at`,不再发生下一次 claim。 ```go // Start() 中调用(不变) if _, err := wp.RequeueStaleJobs(ctx); err != nil { return err } // sweepLoop 中新增定期调用 func (wp *WorkerPool) sweepLoop(ctx context.Context) { // ... sweepDueJobs ... // 新增:定期回收 stale running jobs if _, err := wp.RequeueStaleJobs(ctx); err != nil { applogger.L().Warnf("requeue stale jobs failed: %v", err) } } ``` **影响分析**:`RequeueStaleJobs` 方法签名不变。只有未耗尽的 stale job 会被下一轮 sweep 重新 XADD 到 Redis;耗尽的 job 保持 `dead` 终态。 ### 3.6 Start/Stop 变更 **Start() 新增**: - 启动 sweep goroutine - 确保 consumer group 已创建(`XGROUP CREATE`,如果不存在) ```go func (wp *WorkerPool) Start() error { if wp.db == nil { return nil } wp.mu.Lock() if wp.cancel != nil { wp.mu.Unlock() return nil } wp.ctx, wp.cancel = context.WithCancel(context.Background()) ctx := wp.ctx workerCount := wp.workerCount wp.mu.Unlock() if _, err := wp.RequeueStaleJobs(ctx); err != nil { return err } // 新增:确保 Redis consumer group 存在 if wp.rdb != nil { wp.ensureConsumerGroups(ctx) } for i := 0; i < workerCount; i++ { wp.wg.Add(1) go wp.run(ctx) } // 新增:启动 sweep goroutine if wp.rdb != nil { wp.wg.Add(1) go wp.sweepLoop(ctx) } return nil } func (wp *WorkerPool) ensureConsumerGroups(ctx context.Context) { for _, queue := range wp.allQueues() { streamKey := wp.streamKeyFor(queue) // XGROUP CREATE 幂等:如果 group 已存在则忽略 error err := wp.rdb.XGroupCreate(ctx, streamKey, wp.consumerGroup, "$").Err() if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") { applogger.L().Warnf("XGroupCreate for stream %s group %s: %v", streamKey, wp.consumerGroup, err) } } } ``` **Stop() 不变**:`cancel()` + `wg.Wait()` 已经能正确停止所有 goroutine(包括新增的 sweep)。 ### 3.7 Fallback 模式(Redis 不可用) 当 `rdb == nil`(Redis 未配置或初始化失败)时,worker 自动回退到当前 DB 轮询模式: ```go func (wp *WorkerPool) run(ctx context.Context) { defer wp.wg.Done() if wp.rdb == nil { wp.runDBPollLoop(ctx) // 现有 run() 逻辑,保留 return } // ... Redis XREADGROUP 逻辑 ... } ``` 这保证: - SQLite 测试模式(不依赖 Redis)继续工作 - Redis 故障时 worker 不中断 - 渐进迁移:可以先在部分实例上启用 Redis 队列 ### 3.8 Concurrency 配置修复 **当前 bug**:`Bootstrap()` 使用 `worker.NewWorkerPool(db)`,`cfg.Worker.Concurrency` 未生效。 **修复**: ```go // internal/app/bootstrap.go:109 // 变更前: workerPool := worker.NewWorkerPool(db) // 变更后: workerPool := worker.NewWorkerPoolWithOptions(db, worker.WithRedisClient(rdb), worker.WithWorkerCount(cfg.Worker.Concurrency), ) ``` 注意:`rdb` 在 `Bootstrap()` 中的初始化位于 [bootstrap.go:113](../backend/internal/app/bootstrap.go:113),在 `workerPool` 创建之前。但当前代码先创建 `workerPool`(line 109),后创建 `rdb`(line 113)。需要调整顺序:先创建 `rdb`,再创建 `workerPool`。 ### 3.9 Queue 列表动态化 **当前**:`WorkerPool.queues` 通过 `WithQueues` option 设置,如果未设置则监听所有 queue(`claimNext` 中不追加 queue filter)。 **迁移后**:Redis 模式下必须显式列出要监听的 stream。如果 `queues` 为空,默认监听所有已知 queue。 ```go func (wp *WorkerPool) allQueues() []string { if len(wp.queues) > 0 { return wp.queues } // 默认监听所有 queue return []string{ "default", "high", "medium", "low", "events", "automation", "search", "scheduled_jobs", "deferred", "purgable", } } func (wp *WorkerPool) streamKeys() []string { queues := wp.allQueues() keys := make([]string, 0, len(queues)*2) for _, q := range queues { keys = append(keys, wp.streamKeyFor(q)) keys = append(keys, ">") // > 表示只读取新消息 } return keys } ``` --- ## 4. 配置变更 ### 4.1 config.go 变更 **文件**:`internal/config/config.go` ```go type WorkerConfig struct { Concurrency int `mapstructure:"concurrency"` // 新增 StreamPrefix string `mapstructure:"redis_stream_prefix"` // 默认 "gochat:jobs" ConsumerGroup string `mapstructure:"redis_consumer_group"` // 默认 "gochat-workers" BlockTimeoutS int `mapstructure:"redis_block_timeout_s"` // 默认 5 SweepIntervalS int `mapstructure:"redis_sweep_interval_s"` // 默认 30 } ``` 新增 viper 默认值: ```go v.SetDefault("worker.redis_stream_prefix", "gochat:jobs") v.SetDefault("worker.redis_consumer_group", "gochat-workers") v.SetDefault("worker.redis_block_timeout_s", 5) v.SetDefault("worker.redis_sweep_interval_s", 30) ``` 新增环境变量映射: ```go "GOCHAT_WORKER_REDIS_STREAM_PREFIX": "worker.redis_stream_prefix", "GOCHAT_WORKER_REDIS_CONSUMER_GROUP": "worker.redis_consumer_group", "GOCHAT_WORKER_REDIS_BLOCK_TIMEOUT_S": "worker.redis_block_timeout_s", "GOCHAT_WORKER_REDIS_SWEEP_INTERVAL_S": "worker.redis_sweep_interval_s", ``` ### 4.2 config.dev.yaml 变更 ```yaml worker: concurrency: 5 redis_stream_prefix: "gochat:jobs" redis_consumer_group: "gochat-workers-dev" redis_block_timeout_s: 5 redis_sweep_interval_s: 15 # 开发环境更短,便于测试延迟 job ``` ### 4.3 backend/configs/config.production.yaml 变更 ```yaml worker: concurrency: 10 redis_stream_prefix: "gochat:jobs" redis_consumer_group: "gochat-workers" redis_block_timeout_s: 5 redis_sweep_interval_s: 30 ``` ### 4.4 validator.go 变更 ```go // 新增校验 if cfg.Worker.RedisBlockTimeoutS < 1 { return fmt.Errorf("worker.redis_block_timeout_s must be >= 1") } if cfg.Worker.RedisSweepIntervalS < 1 { return fmt.Errorf("worker.redis_sweep_interval_s must be >= 1") } ``` ### 4.5 hot-reload 支持 `worker.concurrency`、`worker.redis_block_timeout_s`、`worker.redis_sweep_interval_s` 加入 `ReloadableFields`: ```go // config.go ReloadableFields "worker.redis_block_timeout_s", "worker.redis_sweep_interval_s", ``` 注意:`worker.concurrency` 改变需要重启 worker goroutine 才能生效,标记为需要重启(已在不可变字段列表中)。`stream_prefix` 和 `consumer_group` 改变也需要重启。`block_timeout_s` 和 `sweep_interval_s` 可以热更新(下一轮循环自动使用新值)。 --- ## 5. 迁移后受影响功能的验证计划 ### 5.1 核心路径验证 | 验证项 | 验证方法 | 预期结果 | 影响的 job type | |--------|----------|----------|-----------------| | 基本入队+执行 | `go test ./internal/worker/ -run TestWorkerPoolProcessOne` | job 从 queued → running → completed | 任意 | | 幂等去重 | `go test ./internal/worker/ -run TestWorkerPoolEnqueueStoresPayloadAndIdempotency` | 重复 idempotency_key 返回同一 job | 19 种带 key 的 job | | 失败重试+死信 | `go test ./internal/worker/ -run TestWorkerPoolRetriesThenDeadLetters` | attempts 达到 max 后 dead | 任意 | | 延迟调度 | `go test ./internal/worker/ -run TestWorkerPoolRespectsScheduleAndQueues` | scheduled_at 未到时不执行 | `scheduled:trigger_items` 等 5 种 | | Stale 回收 | `go test ./internal/worker/ -run TestWorkerPoolRequeuesStaleRunningJobs` | running 超时 → retrying | 任意 | | Start/Stop 循环 | `go test ./internal/worker/ -run TestWorkerPoolStartAndStop` | goroutine 正确启动和停止 | 任意 | ### 5.2 Redis 队列路径验证(新增测试) | 验证项 | 验证方法 | 预期结果 | |--------|----------|----------| | XADD + XREADGROUP 端到端 | miniredis 创建 stream,Enqueue 后 XREADGROUP 拿到 job_id,DB claim + perform 后 completed | job 状态正确流转 | | Redis 投递失败降级 | mock Redis 返回 error,Enqueue 仍成功返回 job,sweep 后重新投递 | job 不丢失 | | Redis 不可用 fallback | rdb=nil,回退 DB 轮询,ProcessOne 正常工作 | 兼容现有测试 | | 多 consumer 竞争 | 两个 WorkerPool 实例同一 consumer group,同一 job 只被一个执行 | DB claim 检查生效 | | Sweep 补偿 | Enqueue 时 XAdd 失败,sweep 后 XAdd 成功,job 被执行 | 补偿机制有效 | | 延迟 job sweep | Enqueue 时 scheduled_at 未来时间,sweep 到期后投递 | 延迟 job 被投递 | | Consumer group 创建 | 首次 Start 时 XGROUP CREATE,已存在时不报错 | 幂等创建 | ### 5.3 各功能模块验证 #### 5.3.1 消息发送 | 验证项 | 方法 | 预期 | |--------|------|------| | `message:send_reply` 正常投递 | 触发对话回复,检查 channel provider 被调用 | 消息通过 Redis 队列投递到 worker 执行 | | 高优先级 queue 先消费 | 同时入队 `high` 和 `low` job,检查执行顺序 | `high` queue 先被消费 | #### 5.3.2 事件分发 | 验证项 | 方法 | 预期 | |--------|------|------| | `event:dispatch_async` 投递 | 触发 channel event,检查 listener 被调用 | 事件通过 Redis 队列分发 | | `event:listener_dispatch` 投递 | dispatch.EventDispatcher 异步分发 | 正确投递和执行 | #### 5.3.3 自动化规则 | 验证项 | 方法 | 预期 | |--------|------|------| | Automation rule 触发 action delivery | 触发规则条件,检查 email/webhook/transcript delivery job 入队 | job 入 Redis Stream 并被执行 | | Macro 执行 | 调用 macro,检查 `automation:macro_execution` job | 正常执行 | | CSAT survey 发送 | resolve 对话,检查 `csat:survey_send` job | 正常执行,幂等不重复 | #### 5.3.4 定时任务 | 验证项 | 方法 | 预期 | |--------|------|------| | `scheduled:trigger_items` 1h 周期 | 检查自我重新入队 | 延迟 1h 后被 sweep 投递 | | `sla:trigger_accounts` 5min 周期 | 检查自我重新入队 | 延迟 5min 后被 sweep 投递 | | `captain:documents_schedule_syncs` 24h 周期 | 检查自我重新入队 | 延迟 24h 后被 sweep 投递 | | `conversation:update_message_status` 延迟 | 检查延迟投递 | 按指定时间延迟 | | `captain:conversation_response_builder` 附件等待 | 检查 1–5s 延迟 | 按计算时间延迟 | #### 5.3.5 Webhook 入站消息 | 验证项 | 方法 | 预期 | |--------|------|------| | `webhook:incoming_message_persist` | 模拟 FB/WA/TG webhook 回调 | 消息持久化 job 通过 Redis 队列执行 | | `webhook:message_status_update` | 模拟状态回调 | 状态更新 job 正常执行 | | `webhook:contact_messages_status_update` | 模拟联系人状态回调 | 正常执行 | #### 5.3.6 搜索索引 | 验证项 | 方法 | 预期 | |--------|------|------| | `search:index` 正常投递 | 触发消息/对话/联系人变更 | 搜索索引同步 job 通过 Redis 队列执行 | #### 5.3.7 Captain AI | 验证项 | 方法 | 预期 | |--------|------|------| | `captain:copilot_response` | 发送 copilot 消息 | 正常执行 | | `captain:conversation_response_builder` | 入站消息触发 | 延迟后正确执行 | | `captain:document_sync/crawl` | 文档同步操作 | 正常执行 | | `captain:documents_schedule_syncs` | 24h 定时同步 | 延迟后正确执行 | #### 5.3.8 Contact Import/Export | 验证项 | 方法 | 预期 | |--------|------|------| | `contact:import` | 导入 CSV | 正常执行 | | `contact:export` | 导出联系人 | 正常执行 | #### 5.3.9 Conversation Maintenance | 验证项 | 方法 | 预期 | |--------|------|------| | `conversation:reopen_snoozed` | snooze 到期 | 正常 reopen | | `conversation:resolution` | auto-resolve 到期 | 正常 resolve | | `conversation:bulk_action` | 批量操作 | 正常执行 | | `contact:bulk_action` | 批量操作 | 正常执行 | | `conversation:delete_object` | 删除对话 | 正常执行 | #### 5.3.10 SLA 处理 | 验证项 | 方法 | 预期 | |--------|------|------| | `sla:trigger_accounts` | 5min 周期触发 | 延迟后正确执行 | | `sla:process_account` | account 级处理 | 正常执行 | | `sla:process_applied` | applied SLA 处理 | 正常执行 | #### 5.3.11 Reporting | 验证项 | 方法 | 预期 | |--------|------|------| | `reporting:rollup_day` | 触发日汇总 | 正常执行 | ### 5.4 部署验证 | 验证项 | 方法 | 预期 | |--------|------|------| | Quickstart compose | `cd deploy/quickstart && docker compose up -d --build` | 单实例 API+worker 正常运行 | | Prod compose worker 分离 | `docker-compose.prod.yml` worker service | `serve --worker-only` 正常工作(需实现) | | CI postgres + redis | `go test -race ./internal/... ./pkg/... ./cmd/...` | 全部通过 | | CI sqlite (无 redis) | `GOCHAT_TEST_DB=sqlite go test ./internal/... ./pkg/... ./cmd/...` | fallback 模式全部通过 | ### 5.5 回归测试清单 迁移后需要全部通过的现有测试文件(25 个): | 测试文件 | 依赖的 worker 行为 | |----------|-------------------| | `internal/worker/worker_test.go` | Enqueue、ProcessOne、retry、dead-letter、schedule、stale、start/stop | | `internal/dispatch/dispatcher_worker_test.go` | EventDispatcher async job | | `internal/channel/dispatcher_worker_test.go` | channel.Dispatcher async job | | `internal/automation/action_service_test.go` | Action delivery job | | `internal/automation/macro_service_test.go` | Macro execution job | | `internal/automation/csat_survey_listener_test.go` | CSAT survey send job | | `internal/service/conversation_maintenance_worker_test.go` | 全部 maintenance job | | `internal/service/sla_worker_test.go` | SLA processing job | | `internal/service/csat_template_worker_test.go` | CSAT template create job | | `internal/service/copilot_response_worker_test.go` | Copilot response job | | `internal/service/captain_conversation_worker_test.go` | Captain conversation builder job | | `internal/service/captain_document_service_test.go` | Captain document sync/crawl job | | `internal/service/search_indexer_worker_test.go` | Search index job | | `internal/service/message_delivery_worker_test.go` | Message send reply job | | `internal/service/contact_service_g3_test.go` | Contact import/export job | | `internal/service/conversation_service_test.go` | Conversation delete job | | `internal/service/inbox_service_test.go` | Inbox template sync job | | `internal/service/widget_service_test.go` | Webwidget triggered job | | `internal/service/article_service_test.go` | Article translate job | | `internal/service/message_service_test.go` | Message service job | | `internal/service/analytics_p513_test.go` | Analytics job | | `internal/handler/webhook/webhook_lookup_test.go` | Webhook job | | `internal/handler/api/v1/inbox_handler_test.go` | Inbox handler job | | `internal/handler/api/v1/bulk_action_handler_test.go` | Bulk action job | | `internal/handler/api/v1/article_handler_test.go` | Article handler job | --- ## 6. 实施步骤 ### Phase 1:核心改造(不改变对外接口) | 步骤 | 文件 | 改动 | |------|------|------| | 1.1 | `internal/worker/worker.go` | WorkerPool 结构体加 Redis 字段;新增 `WithRedisClient` 等 Option | | 1.2 | `internal/worker/worker.go` | `Enqueue` 末尾加 `pushToRedis` | | 1.3 | `internal/worker/worker.go` | `run` 改为 `XREADGROUP BLOCK`;保留 `runDBPollLoop` fallback | | 1.4 | `internal/worker/worker.go` | 新增 `processRedisMessage`、`ackRedis`、`pushToRedis`、`streamKeyFor`、`streamKeys`、`allQueues` | | 1.5 | `internal/worker/worker.go` | 新增 `sweepLoop`、`sweepDueJobs` | | 1.6 | `internal/worker/worker.go` | `Start` 加 `ensureConsumerGroups` + sweep goroutine | | 1.7 | `internal/worker/worker.go` | `ProcessOne`、`claimNext`、`perform`、`fail` 保留不变 | ### Phase 2:配置与接入 | 步骤 | 文件 | 改动 | |------|------|------| | 2.1 | `internal/config/config.go` | WorkerConfig 加 4 个 Redis 字段 | | 2.2 | `internal/config/config.go` | viper SetDefault + 环境变量映射 | | 2.3 | `internal/config/config.go` | applyZeroDefaults 补充新字段 | | 2.4 | `internal/config/validator.go` | 新增字段校验 | | 2.5 | `internal/app/bootstrap.go` | 调整 rdb 创建顺序(移到 workerPool 之前),改用 `NewWorkerPoolWithOptions` | | 2.6 | `configs/config.dev.yaml` | 新增 worker Redis 配置 | | 2.7 | `backend/configs/config.production.yaml` | 新增 worker Redis 配置 | ### Phase 3:测试 | 步骤 | 文件 | 改动 | |------|------|------| | 3.1 | `internal/worker/worker_test.go` | 新增 miniredis 测试:Redis 端到端、竞争、sweep、fallback | | 3.2 | `internal/worker/worker_test.go` | 现有测试保持:rdb=nil 时 fallback 到 DB 轮询 | | 3.3 | 全部 worker 相关测试 | `GOCHAT_TEST_DB=sqlite` 跑一遍确认 fallback 模式通过 | ### Phase 4:独立 worker 进程(可选,后续迭代) | 步骤 | 文件 | 改动 | |------|------|------| | 4.1 | `cmd/gochat/main.go` | 新增 `serve --worker-only` 子命令 | | 4.2 | `internal/app/bootstrap.go` | 拆出 `BootstrapWorker()` 函数,只初始化 DB + Redis + worker 相关依赖 | | 4.3 | `deploy/docker/Dockerfile` | 更新注释,`gochat-worker` 二进制对应 `serve --worker-only` | | 4.4 | `deploy/docker/docker-compose.prod.yml` | worker service command 改为实际可用命令 | --- ## 7. 风险与缓解 | 风险 | 影响 | 缓解 | |------|------|------| | Redis 宕机 | worker 无法拉取新 job | 自动 fallback 到 DB 轮询;已有 job 在 DB 中不丢失 | | Redis Stream 消息积压 | 内存增长 | 配置 `MAXLEN` 限制 stream 长度;sweep + XACK 清理已处理消息 | | Sweep 延迟过高 | 延迟 job 执行晚于预期 | 可配置 `sweep_interval`(默认 30s,开发环境 15s) | | 多实例 worker 重复执行 | 同一 job 被执行两次 | DB claim 检查(`WHERE status IN ... AND ...` + `RowsAffected` 检查)保证只有一个成功 | | Redis Stream 不保证顺序 | 同一 queue 内 job 执行顺序不保证 | 当前 DB 轮询也不保证多 worker 间的顺序;`priority` 改为 queue 级别优先级 | | 迁移期间混用 | 部分 instance 用 DB 轮询,部分用 Redis | DB claim 检查兼容两种 consumer;Redis XADD + DB INSERT 同时存在不会冲突 | --- ## 8. 不改动的部分 以下内容在迁移中**完全不改动**: - `BackgroundJob` model(字段、索引) - `background_jobs` 表 schema 和 migration - 所有 `Register*Jobs` 函数(`RegisterConversationMaintenanceJobs`、`RegisterActionDeliveryJobs` 等) - 所有 `Enqueue*` 辅助函数(`EnqueueSendReply`、`EnqueueScheduledItemsTrigger` 等) - 所有 `SetWorkerPool` 调用 - 所有 `JobHandler` 实现 - 所有 job type 常量 - 所有 queue 名称 - `EnqueueOption` / `WithQueue` / `WithScheduledAt` / `WithMaxAttempts` / `WithPriority` / `WithIdempotencyKey` 函数签名 - `perform()` 和 `fail()` 方法实现 - `RequeueStaleJobs()` 方法签名和 DB 操作 - `defaultBackoff()` 函数 - `marshalPayload()` 函数 - Redis Pub/Sub 广播(`internal/ws/broadcast.go`) - Watermill EventBus(`internal/pubsub/event_bus.go`) - Watermill RedisPubSub(`internal/pubsub/redis_pubsub.go`) - NotificationDeliveryService(`internal/service/notification_delivery_service.go`) --- ## 9. 文件改动清单 | 文件 | 类型 | 改动范围 | |------|------|----------| | `internal/worker/worker.go` | 核心改动 | 结构体字段、Enqueue、run、Start、新增 processRedisMessage/sweep | | `internal/worker/worker_test.go` | 新增测试 | miniredis Redis 队列路径测试 | | `internal/config/config.go` | 配置扩展 | WorkerConfig 结构体、defaults、env 映射 | | `internal/config/validator.go` | 校验 | 新字段校验 | | `internal/app/bootstrap.go` | 接入 | rdb 顺序调整、NewWorkerPoolWithOptions | | `configs/config.dev.yaml` | 配置 | worker Redis 参数 | | `backend/configs/config.production.yaml` | 配置 | worker Redis 参数 | | `deploy/docker/Dockerfile` | 注释更新 | gochat-worker 说明(Phase 4) | | `deploy/docker/docker-compose.prod.yml` | 部署 | worker service command(Phase 4) |