Files
gochat/docs/architecture/worker-redis-queue-migration.md
2026-08-24 17:01:17 +08:00

922 lines
39 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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 延迟 Jobscheduled_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 {
// fallbackRedis 不可用时回退到 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 创建 streamEnqueue 后 XREADGROUP 拿到 job_idDB claim + perform 后 completed | job 状态正确流转 |
| Redis 投递失败降级 | mock Redis 返回 errorEnqueue 仍成功返回 jobsweep 后重新投递 | 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 检查兼容两种 consumerRedis 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 commandPhase 4 |