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

39 KiB
Raw Blame History

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.goWorkerPool 实现,核心行为:

  • 入队Enqueue()BackgroundJob 记录 INSERT 到 PostgreSQL background_jobs
  • 消费:每个 worker goroutine 每 500ms 轮询 DBSELECT ... 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 定义,仅改变 Enqueuerun 的内部实现。

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 keygochat:jobs:{queue}

2.3 延迟 Jobscheduled_at)的使用场景

5 个 job type 使用 WithScheduledAt,需要特殊处理:

  1. scheduled:trigger_items1 小时周期,自我重新入队
  2. sla:trigger_accounts5 分钟周期,自我重新入队
  3. captain:documents_schedule_syncs24 小时周期,自我重新入队
  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 DESCBackgroundJob.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 变更

当前流程

  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 <= nowXADD 到 Redis Stream
  5. 新增:如果 scheduled_at > now,跳过 XADD(由 sweep 机制在到期时投递)
  6. 返回 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* 辅助函数(EnqueueSendReplyEnqueueScheduledItemsTrigger 等)不需要改动,它们调用 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 {
        // 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。

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 配置修复

当前 bugBootstrap() 使用 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),
)

注意:rdbBootstrap() 中的初始化位于 bootstrap.go:113,在 workerPool 创建之前。但当前代码先创建 workerPoolline 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.concurrencyworker.redis_block_timeout_sworker.redis_sweep_interval_s 加入 ReloadableFields

// config.go ReloadableFields
"worker.redis_block_timeout_s",
"worker.redis_sweep_interval_s",

注意:worker.concurrency 改变需要重启 worker goroutine 才能生效,标记为需要重启(已在不可变字段列表中)。stream_prefixconsumer_group 改变也需要重启。block_timeout_ssweep_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 先消费 同时入队 highlow 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 附件等待 检查 15s 延迟 按计算时间延迟

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 新增 processRedisMessageackRedispushToRedisstreamKeyForstreamKeysallQueues
1.5 internal/worker/worker.go 新增 sweepLoopsweepDueJobs
1.6 internal/worker/worker.go StartensureConsumerGroups + sweep goroutine
1.7 internal/worker/worker.go ProcessOneclaimNextperformfail 保留不变

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 函数(RegisterConversationMaintenanceJobsRegisterActionDeliveryJobs 等)
  • 所有 Enqueue* 辅助函数(EnqueueSendReplyEnqueueScheduledItemsTrigger 等)
  • 所有 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 EventBusinternal/pubsub/event_bus.go
  • Watermill RedisPubSubinternal/pubsub/redis_pubsub.go
  • NotificationDeliveryServiceinternal/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