游戏资源派生队列可以丢任务,但不能丢 pending 状态

本文代码来自一个真实 Go 在线内容服务的图片派生流水线;把同一约束迁移到游戏活动纹理、启动器封面或资源预览时,边界不变。若上传请求等待全部转码,解码抖动会进入运营链路;若只把任务塞进内存队列,进程重启又会永久漏图。正确边界是:数据库保存原图和 pending 欠账,队列只负责低成本唤醒。

// ../src/internal/imgproc/pipeline.go: Pipeline
const defaultWorkers = 2
const queueCapacity = 64

type Pipeline struct {
    st      *store.Store
    mode    string
    fsDir   string
    queue   chan int64
    wg      sync.WaitGroup
    cancel  context.CancelFunc
    mu      sync.Mutex
    started bool
    stopped bool
    stopOnce sync.Once
}

func NewPipeline(st *store.Store, mode, fsDir string) *Pipeline {
    return &Pipeline{
        st: st, mode: mode, fsDir: fsDir,
        queue: make(chan int64, queueCapacity),
    }
}

一张权威原始游戏纹理通过持久状态进入双 Worker 派生流水线,瞬时队列溢出不会损坏原图

先记欠账

Pipeline 不是资源所有者。真正需要恢复的是数据库中的原图和 transcode_status=pending。通道只保存 origin_id,即使整个通道消失,数据库仍能回答“还有哪些原图没有形成终态”。

// ../src/internal/store/variants.go: PendingTranscodeImages
func (s *Store) PendingTranscodeImages(ctx context.Context, limit int) ([]int64, error) {
    if limit <= 0 {
        limit = 50
    }
    rows, err := s.DB.QueryContext(ctx,
        `SELECT id FROM images
         WHERE transcode_status='pending'
         ORDER BY id ASC LIMIT ?`, limit)
    if err != nil {
        return nil, err
    }
    defer rows.Close()

    var ids []int64
    for rows.Next() {
        var id int64
        if err := rows.Scan(&id); err != nil {
            return nil, err
        }
        ids = append(ids, id)
    }
    return ids, rows.Err()
}

这里的不变量只有一条最重要:业务入口必须先持久化原图与 pending,再调用 Enqueue。若顺序反过来,worker 可能先读到不存在的原图;若只有队列没有状态表,进程崩溃后无法区分“从未提交”和“等待处理”。

队列只唤醒

上传链路不能被派生任务反压,因此 Enqueue 使用带 defaultselect。满队列丢掉的是一次唤醒,不是原图,也不是工作事实。

// ../src/internal/imgproc/pipeline.go: Enqueue
func (p *Pipeline) Enqueue(originID int64) {
    if originID <= 0 {
        return
    }
    p.mu.Lock()
    stopped := p.stopped
    p.mu.Unlock()
    if stopped {
        return
    }
    defer func() {
        if rec := recover(); rec != nil {
            log.Printf("imgproc: enqueue after stop, drop id=%d", originID)
        }
    }()
    select {
    case p.queue <- originID:
    default:
        log.Printf("imgproc: queue full, drop id=%d", originID)
    }
}

这项取舍适合可重建的派生资源,不适合付费发奖、角色存档或对局结算。那些任务本身就是权威事实,不能把“稍后扫描”当作提交协议。

启动补偿

Start 创建两个 worker 后立即回扫 pending。当前单次上限是 200;它保证小规模欠账能恢复,不证明任意积压都能在启动窗口内清空。

// ../src/internal/imgproc/pipeline.go: Start + scanPending
func (p *Pipeline) Start(ctx context.Context) {
    wctx, cancel := context.WithCancel(ctx)
    p.cancel = cancel
    for i := 0; i < defaultWorkers; i++ {
        p.wg.Add(1)
        go p.worker(wctx, i)
    }
    go p.scanPending(wctx)
}

func (p *Pipeline) scanPending(ctx context.Context) {
    ids, err := p.st.PendingTranscodeImages(ctx, 200)
    if err != nil {
        log.Printf("scan pending: %v", err)
        return
    }
    for _, id := range ids {
        select {
        case <-ctx.Done():
            return
        case p.queue <- id:
        }
    }
}

正常时序是 pending → 入队 → process → variants → done。失败时序是 pending → 队列满或进程退出 → 唤醒丢失 → 下次启动回扫 → process。恢复能力来自状态表,不来自通道容量。

原图 pending 状态、容量 64 的非阻塞队列、双 Worker、启动回扫与实际运行数据

写入终态

Worker 把 panic 和普通错误都收敛为 error;成功才写 doneprocess 发现已有变体会直接跳过,使重复唤醒不会重复生成。

// ../src/internal/imgproc/pipeline.go: worker
storageKey, err := p.process(ctx, id)
if err != nil {
    _ = p.st.SetImageTranscodeStatus(ctx, id, "error", err.Error())
    return
}
_ = p.st.SetImageTranscodeStatus(ctx, id, "done", "")

// ../src/internal/imgproc/pipeline.go: process
if existing, _ := p.st.ListVariantsByOrigin(ctx, originID); len(existing) > 0 {
    return "", nil
}
widths := []int{originW}
for _, w := range []int{800, 1200} {
    if w < originW {
        widths = append(widths, w)
    }
}
for _, w := range widths {
    resized := Resize(src, w)
    data, err := EncodeWebP(resized)
    if err != nil {
        return "", err
    }
    if err := p.persistVariant(ctx, img, "webp", w, resized.Bounds().Dy(), data); err != nil {
        return "", err
    }
}

这里仍有真实失败窗口:若前两个变体已落库、第三个失败,重试会因“已有任意变体”整体跳过,随后 worker 可能写成 done。当前测试没有覆盖部分成功。更严格的实现应按预期宽度集合逐项补齐,或在一个事务中提交完整 manifest。

测试约束

测试直接执行真实 JPEG 编码、SQLite 迁移和变体查询。1600×1000 输入必须得到 1600、1200、800 三档;600×400 小图只能保留 600,不能上采样。

// ../src/internal/imgproc/pipeline_test.go: TestPipelineProcessJPEG
p := NewPipeline(st, "db", "")
if _, err := p.process(context.Background(), img.ID); err != nil {
    t.Fatalf("process: %v", err)
}
vs, err := st.ListVariantsByOrigin(context.Background(), img.ID)
if err != nil {
    t.Fatalf("list: %v", err)
}
if len(vs) != 3 {
    t.Fatalf("want 3 variants, got %d", len(vs))
}

// TestEnqueueDoesNotBlock
p = NewPipeline(nil, "db", "")
for i := 0; i < queueCapacity+50; i++ {
    p.Enqueue(int64(i + 1))
}
$ go test -count=1 -v ./internal/imgproc -run 'Test(...)$'
PASS: 6 / 6, package 1.560s, real 3.10s
TestPipelineProcessJPEG        0.51s
TestPipelineSkipsSmallImage    0.04s
TestEnqueueDoesNotBlock        0.00s
queue full: dropped id=65 ... id=114   # 50 次

$ go test -count=1 -v ./internal/store -run '^TestImageVariants$'
PASS: 1 / 1, package 0.432s, real 0.65s

$ go test -race -count=1 ./internal/imgproc ./internal/store
ok imgproc 15.181s; ok store 4.722s; real 22.11s
exit code: 0

运行数据

只读本地数据库快照进一步说明派生资源的容量属性。5 张原图都处于 done,共产生 13 个变体;原图总计 1,017,921 字节,变体总计 3,687,854 字节,是原图的 3.62 倍。派生不是免费缓存,必须单独规划容量和清理策略。

SELECT transcode_status, count(*) FROM images GROUP BY transcode_status;
-- done | 5

SELECT count(*), sum(size) FROM image_variants;
-- 13 | 3687854

SELECT width, count(*) FROM image_variants GROUP BY width ORDER BY width;
-- 741:1, 800:4, 1200:4, 1473:1, 1600:1, 2320:1, 3076:1

当前设计用 64 个内存槽和两个 worker 换取上传链路稳定,代价是恢复依赖启动回扫、错误项不会自动重试、部分成功缺少 manifest 约束。规模扩大后,先增加 pending 年龄、error 数量和队列丢弃计数;当欠账经常超过 200、单机恢复时间不可接受或派生物需要跨节点生成时,再把扫描改成可租约的数据库任务表。无需先上消息系统:只要权威欠账仍在数据库,队列就可以继续保持简单、可丢和便宜。