内容审核流水线:文本、图片和视频的异步审核架构设计
内容审核流水线:文本、图片和视频的异步审核架构设计
一、多模态内容爆发下的审核架构困局
UGC平台日均内容发布量突破百万级别后,同步审核模式的瓶颈会直接暴露在生产监控面板上。一次同步审核请求串行调用文本敏感词匹配、图片鉴黄模型和视频关键帧抽取,端到端耗时落在 800ms 到 2s 之间。当上游业务方的 HTTP 超时设为 1s 时,超过一半的审核请求直接超时丢弃——这意味着该被拦截的内容漏过去了。
问题不在于单一模型太慢,而在于耦合。文本、图片、视频三种模态的审核逻辑被塞进同一个请求-响应周期里,慢的环节拖死快的环节。更致命的是,视频审核需要先做关键帧抽取再做逐帧推理,天然无法在 1s 内完成。
基础设施不需要漂亮话,需要的是把三种模态的审核拆成独立流水线、各自异步执行的架构。
二、流水线解耦:消息队列驱动的多模态异步架构
核心思路是把"审核"这件事从同步接口调用变成任务提交-结果回调的异步模型。业务方发布内容后,审核服务只做两件事:接收内容元数据、将审核任务投递到对应的消息队列,然后立即返回任务 ID。
每种模态对应一条独立的消费流水线。
文本审核流水线按优先级分为两级:第一级是 AC 自动机做敏感词匹配,耗时通常 5ms 以内;第二级是 NLP 模型做语义级别的违规识别,耗时 50ms-200ms。两级之间通过同一个队列的不同 Consumer Group 串联——第一级判定违规的直接终止流水线,只有疑似或通过的内容才进入第二级。
图片审核的挑战在于图片本身不跟请求体一起提交。业务方传的是 CDN URL,Worker 需要先下载图片再做推理。这里不能阻塞在下载上——单 Worker 用同步 HTTP GET 下载图片的 P99 延迟可能到 2s。做法是每个 Worker 内部用 goroutine 池并发下载多张图片,下载超时设为 500ms,超时直接标记为"审核失败-下载超时"。
视频审核是三层漏斗:关键帧抽取 → 图片审核(复用图片流水线)→ 音频转文字审核。关键帧抽取用 FFmpeg 按 I 帧间隔抓帧,10 分钟视频约产生 600 帧,抽帧后逐帧投递到图片审核队列,最后汇总所有帧的审核结果取最严判定。
三、Go 生产级任务路由与 Worker Pool 实现
任务路由的核心是把请求中的内容类型映射到对应的消息队列 Topic 上。下面是一个简化但可投产的实现思路。
type ModerationTask struct { TaskID string `json:"task_id"` BizID string `json:"biz_id"` MediaType string `json:"media_type"` // text/image/video Content string `json:"content"` // text or URL Callback string `json:"callback"` CreatedAt time.Time `json:"created_at"` } // RouteTask 根据媒体类型将任务路由到对应队列 func (r *TaskRouter) RouteTask(ctx context.Context, task ModerationTask) error { var topic string switch task.MediaType { case "text": topic = "moderation.text" case "image": topic = "moderation.image" case "video": topic = "moderation.video" default: return fmt.Errorf("unsupported media type: %s", task.MediaType) } payload, err := json.Marshal(task) if err != nil { return fmt.Errorf("marshal task: %w", err) } // RabbitMQ 投递,消息持久化 + 手动确认 return r.ch.PublishWithContext(ctx, "moderation.exchange", // exchange topic, // routing key true, // mandatory false, // immediate amqp.Publishing{ DeliveryMode: amqp.Persistent, ContentType: "application/json", Body: payload, MessageId: task.TaskID, }, ) }Worker Pool 侧的关键设计是并发度与 QPS 的匹配。假设单 Worker 处理一条文本审核耗时 10ms,期望吞吐 1000 QPS,理论上 10 个 goroutine 就够了。但实际中网络抖动、GC 停顿会让 P99 膨胀到 30ms,所以需要按 P99 延迟来算并发数:concurrency = target_qps * p99_latency,即1000 * 0.03 = 30。预留 1.5 倍余量,实际开 45 个 goroutine。
type WorkerPool struct { tasks <-chan ModerationTask handler TaskHandler workers int } func (wp *WorkerPool) Start(ctx context.Context) { for i := 0; i < wp.workers; i++ { go func(workerID int) { for { select { case <-ctx.Done(): return case task, ok := <-wp.tasks: if !ok { return } // 每条消息独立超时控制,不因单条慢消息阻塞整个 Worker taskCtx, cancel := context.WithTimeout(ctx, 30*time.Second) if err := wp.handler.Handle(taskCtx, task); err != nil { log.Printf("worker %d handle task %s: %v", workerID, task.TaskID, err) } cancel() } } }(i) } }四、异步架构的代价:最终一致性与可观测性
异步化不是免费的。同步模式下,业务方发请求后 1s 内拿到审核结果,直接决定内容是否放行。异步模式把"放行"和"审核完成"解耦了,引入了一个时间窗口:内容已发布但审核还没跑完。
解决方式有两种。第一种是先发后审,内容发布后先进入"审核中"状态,用户自己可见但不对公展示,审核通过后自动切换为公开。第二种是先审后发,发布接口同步等待,但等待的只是"任务已提交"而非"审核已完成",发布后内容立即进入"审核中"。
这不是技术选择问题,是产品策略问题。工程上需要保证的是:无论选哪种模式,审核结果的回调必须可靠。
可观测性方面,每个审核任务需要贯穿全链路的 trace_id。从任务路由 → 队列投递 → Worker 消费 → 模型推理 → 结果回调,每一步都打上 span。当某条内容的审核延迟超过 SLO 时,能直接定位到是队列积压、模型推理慢还是回调超时。
另外一条硬规:审核任务必须有死信队列。消息被 Nack 或消费超时的任务不能直接丢弃,必须进入死信队列等待人工介入或定时重投。内容审核的漏审成本远高于重复审核成本。
五、总结
内容审核的异步流水线架构,本质是把"不可预测的模态差异"通过消息队列转化为"可独立扩缩的消费单元"。三个要点:
- 按模态拆分队列:文本、图片、视频各自独立,避免慢模态拖死快模态。
- Worker 并发度按 P99 算:不要用平均延迟算并发数,用 P99 峰值加 1.5 倍余量。
- 死信队列不可省略:审核漏过一条违规内容的代价比重复审核一条正常内容高几个数量级。
异步化引入的最终一致性问题可以通过"审核中"状态屏蔽来解决,关键在于回调可靠性和全链路 trace。
