39 KiB
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 到 PostgreSQLbackground_jobs表 - 消费:每个 worker goroutine 每 500ms 轮询 DB,
SELECT ... FOR UPDATE SKIP LOCKEDclaim 最早到期的 job - 执行:
perform()调用注册的JobHandler,完成后 UPDATE 状态为completed/retrying/dead - 补偿:
RequeueStaleJobs()扫描超过staleLockTimeout(默认 15 分钟)的runningjob,重置为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 当前架构的问题
- 轮询开销:500ms 固定轮询,空闲时持续消耗 DB 连接和 CPU
- 延迟下限:job 从入队到执行至少 0–500ms 延迟
- DB 连接占用:每个 worker goroutine 持续占用一个 DB 连接做轮询
- concurrency 配置未生效:
Bootstrap()使用NewWorkerPool(db)而非NewWorkerPoolWithOptions(db, WithWorkerCount(n)),实际只有 1 个 goroutine - 基础设施不统一: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,需要特殊处理:
scheduled:trigger_items:1 小时周期,自我重新入队sla:trigger_accounts:5 分钟周期,自我重新入队captain:documents_schedule_syncs:24 小时周期,自我重新入队conversation:update_message_status:可选延迟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
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 函数:
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 变更
当前流程:
- 构建
BackgroundJob结构 - 处理 idempotency key 去重(DB 查询)
db.Create(job)插入 DB- 返回 job
迁移后流程:
- 构建
BackgroundJob结构(不变) - 处理 idempotency key 去重(DB 查询,不变)
db.Create(job)插入 DB(不变)- 新增:如果
scheduled_at <= now,XADD到 Redis Stream - 新增:如果
scheduled_at > now,跳过 XADD(由 sweep 机制在到期时投递) - 返回 job(不变)
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 阻塞拉取
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。
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。
// 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,如果不存在)
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 轮询模式:
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 未生效。
修复:
// 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,在 workerPool 创建之前。但当前代码先创建 workerPool(line 109),后创建 rdb(line 113)。需要调整顺序:先创建 rdb,再创建 workerPool。
3.9 Queue 列表动态化
当前:WorkerPool.queues 通过 WithQueues option 设置,如果未设置则监听所有 queue(claimNext 中不追加 queue filter)。
迁移后:Redis 模式下必须显式列出要监听的 stream。如果 queues 为空,默认监听所有已知 queue。
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
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 默认值:
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)
新增环境变量映射:
"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 变更
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 变更
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 变更
// 新增校验
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:
// 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. 不改动的部分
以下内容在迁移中完全不改动:
BackgroundJobmodel(字段、索引)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) |