AI 剪辑任务异步管道(MQ 语义模拟)
难度:⭐⭐⭐ 中等偏难
考点
- 生产者 / 消费者模型、任务队列
- 多 worker 并发消费(模拟 Kafka consumer group:一个任务只被一个 worker 消费)
- at-least-once:失败重试(模拟"处理成功才提交 offset")
- 优雅关闭:Close 后拒绝新任务、等待已提交任务处理完
题目描述
用内存队列模拟文档 6.5 的「AI 剪辑任务」异步链路:提交(Producer)→ 队列(Kafka topic)→ 多 worker 消费(Consumer group)→ 处理结果事件(clip-result topic)。
实现 Pipeline:
NewPipeline(ctx, workers, maxRetry, process):启动workers个 worker 并发消费Submit(ctx, task):入队;管道已Close返回ErrClosed;ctx取消返回ctx.Err()- 每个任务最多尝试
maxRetry+1次(首次 + maxRetry 次重试),全部失败则Result.Err非空(at-least-once 语义) Results():返回处理结果事件通道(模拟 clip-result topic)Close():停止接收新任务,等待已提交任务全部处理完,然后关闭Results()(幂等)
函数签名
go
type Task struct {
ID string
Payload string
}
type Result struct {
Task Task
Err error
}
var ErrClosed = errors.New("pipeline closed")
type Pipeline struct{ /* 自行设计 */ }
func NewPipeline(ctx context.Context, workers, maxRetry int, process func(ctx context.Context, t Task) error) *Pipeline
func (p *Pipeline) Submit(ctx context.Context, t Task) error
func (p *Pipeline) Results() <-chan Result
func (p *Pipeline) Close()提示
- 队列用 buffered channel;worker 用
for t := range jobs,Close关闭 jobs 后 worker 排空队列退出 Close里wg.Wait()之后再close(results)(否则发送端 panic)Submit用select同时监听"已关闭"信号与ctx.Done()- 思考:真实 Kafka 中"手动提交 offset"如何保证 at-least-once?(处理成功才提交,崩溃后从未提交处重放,代价是可能重复处理 → 消费端要幂等)
- 注意:
Close与Submit并发调用需要外部加锁(延伸思考,本练习不要求)
验收
- [ ] 6 任务 3 worker:6 条结果事件,全部成功
- [ ] 失败 2 次后成功的任务:
process共被调用 3 次,结果无错误 - [ ] 永远失败的任务:结果带错误(重试耗尽)
- [ ]
Close后Submit返回ErrClosed
参考答案(Go)
点击展开参考答案
go
//go:build ignore
package answer
import (
"context"
"errors"
"sync"
)
type Task struct {
ID string
Payload string
}
type Result struct {
Task Task
Err error
}
var ErrClosed = errors.New("pipeline closed")
// Pipeline 参考答案
type Pipeline struct {
ctx context.Context
jobs chan Task
results chan Result
closed chan struct{}
closeOne sync.Once
wg sync.WaitGroup
process func(ctx context.Context, t Task) error
maxRetry int
}
func NewPipeline(ctx context.Context, workers, maxRetry int, process func(ctx context.Context, t Task) error) *Pipeline {
if workers <= 0 {
workers = 1
}
p := &Pipeline{
ctx: ctx,
jobs: make(chan Task, 64),
results: make(chan Result, 64),
closed: make(chan struct{}),
process: process,
maxRetry: maxRetry,
}
for i := 0; i < workers; i++ {
p.wg.Add(1)
go func() {
defer p.wg.Done()
for t := range p.jobs {
p.handle(t)
}
}()
}
return p
}
func (p *Pipeline) Submit(ctx context.Context, t Task) error {
select {
case <-p.closed:
return ErrClosed
default:
}
if ctx.Err() != nil {
return ctx.Err()
}
select {
case p.jobs <- t:
return nil
case <-p.closed:
return ErrClosed
case <-ctx.Done():
return ctx.Err()
}
}
func (p *Pipeline) Results() <-chan Result {
return p.results
}
func (p *Pipeline) Close() {
p.closeOne.Do(func() {
close(p.closed)
close(p.jobs)
p.wg.Wait()
close(p.results)
})
}
func (p *Pipeline) handle(t Task) {
var err error
for attempt := 0; attempt <= p.maxRetry; attempt++ {
if err = p.process(p.ctx, t); err == nil {
break
}
}
p.results <- Result{Task: t, Err: err}
}