ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

文生视频异步任务网关治理实战:状态机、幂等与可靠回调

文生视频异步任务网关治理实战:状态机、幂等与可靠回调 文生视频这波热度有多高不用我多说但真正把这类业务从demo推到线上稳定跑的人大概都绕不开一件事异步任务怎么治理。视频生成不是普通接口调用一个prompt丢进去显卡要算几十秒甚至几分钟整个交互模型天然就是异步的。文生视频和发邮件、生成图片这种异步任务还有一个本质区别单任务耗时长、成本高、失败容忍度低。排队中的请求能不能及时处理GPU资源够不够任务跑挂了用户会不会原地爆炸这些问题几乎全压在后端链路上而首当其冲的就是任务网关层。这篇文章就聊我做文生视频业务时在异步任务网关层治理上踩过的坑和最终的解决方案。里面涉及的选型思路、状态机设计、代码片段、排查手段都是线上验证过的。适合正在做或准备做AIGC视频类服务的后端同学参考尤其是任务链路比较复杂、需要自己管理任务状态和回调通知的团队。1. 为什么文生视频业务必须认真治理异步任务网关层1.1 从一次线上事故说起回调丢失带来的连环崩溃先讲一个我们真实遇到的事故。上线初期我们的架构非常简单客户端调用生成接口后端创建任务丢给worker生成worker完成后再把视频URL回调给客户端。问题就出在这个回调上。某天晚上高峰回调服务因为一次发版导致的连接误关闭大量回调消息直接丢掉一批任务其实已经生成成功了但客户端永远收不到通知用户在前端看到的是生成中状态一直在转圈。比较尴尬的是任务生成的视频其实都躺在对象存储里客户就是拿不到。这个事故的直接原因是回调不可靠但往深了挖根子在网关层设计。我们没有为回调设计持久化、重试和兜底查询整个链路是尽力而为而不是保证送达。这也是我后来把异步任务网关层单独拎出来治理的原因文生视频这种高成本、长耗时的任务网关层不是简单的转发代理它其实是整个任务生命周期的状态中枢。1.2 异步任务在文生视频架构中的位置从整体架构来看一条文生视频的完整链路大致是客户端 - 网关层任务提交/查询/回调通知 - 队列 - workerGPU推理 - 对象存储 - 状态持久化网关层在这里承担了三类工作。第一承接所有任务请求做参数校验、鉴权、限流第二维护任务状态流转让用户随时能问到我的视频生成到哪一步了第三负责结果通知通过轮询或回调把最终结果送达客户端。说白了网关层就是整个异步任务的交通枢纽。队列保证worker不会被请求打爆状态表让任务有迹可循回调让用户不用一直傻等。这三个环节里任何一个设计得不够严谨后面一定会出问题。1.3 网关层治理要解决的四个核心问题我做了这段时间之后把异步任务网关层的核心问题收敛成四个可追踪每个任务必须能查到全生命周期的状态包括排队中、生成中、成功、失败、超时、取消还得看到每个状态的耗时分布否则出问题完全没法定位。可限流文生视频的GPU资源是稀缺且贵的网关层必须从入口就控制并发和速率否则后端会把显卡打爆其他人全部排队。可送达生成结果无论通过回调还是轮询必须可靠到达客户端手里不允许视频生成了但用户不知道这种埋雷场景。可自愈任务处理过程中难免有网络抖动、worker异常、服务重启网关层要能自动重试、超时置位、补偿修复而不是靠人肉半夜爬起来改数据库。这四个问题其实就是我后面所有设计的出发点。先把这个想清楚再去选型就不会被各种框架带偏。2. 任务状态机与数据模型设计2.1 状态机设计不要只有成功和失败我见过不少团队做异步任务时状态表里就三个值pending、success、fail。短任务这么搞勉强能行但文生视频绝对不行。一个视频任务可能要跑几十秒到几分钟用户需要知道到底是在排队还是已经在生成否则体验非常差。我们最终敲定的状态机是这样的CREATED提交成功任务已经落库还没有入队QUEUED任务已进入消息队列等待worker消费PROCESSINGworker已拉取任务正在生成视频SUCCEEDED视频生成成功结果地址已回写FAILED生成失败记录错误码和原因TIMEOUT任务超时排队超时或处理超时由调度器扫描置位CANCELED用户主动取消每个状态还包括辅助字段比如重试次数、最后更新时间等。状态流转规则我用一个小表列出来方便大家对照当前状态可流转到触发方CREATEDQUEUED / FAILED / CANCELED入队成功发消息入队失败用户取消QUEUEDPROCESSING / FAILED / TIMEOUTworker拉取worker异常调度器超时扫描PROCESSINGSUCCEEDED / FAILED / TIMEOUTworker回写worker报错调度器超时扫描SUCCEEDED-终态-FAILEDQUEUED自动重试 / -重试调度人工处理TIMEOUTCANCELED / FAILED用户取消重试失败这里踩过的一个坑是一开始我们没区分CREATED和QUEUED认为提交即入队。后来发现消息发到队列也可能失败比如队列抖动如果状态表里只有提交成功用户看到的状态就会误导人。所以建议宁可多两个状态也不要为了简单丢失中间状态。2.2 任务ID和幂等键一个单独的全局ID还不够任务ID本身很简单我们用雪花算法生成64位ID数据库主键用bigint对外返回字符串。真正容易踩坑的是幂等键。文生视频这种高成本业务一定要支持幂等提交。因为客户端的网络环境非常不稳定一个请求超时后客户端会自动重试如果后端没有幂等处理同一个prompt可能被提交进去两次等于GPU白白算了两遍成本直接翻倍。更麻烦的是用户下一次刷新页面还能看到两次记录体验也很差。我们的做法是客户端在提交请求时根据用户ID prompt内容 参数hash生成一个idempotency_key也允许客户端自己传一个UUID服务端在创建任务前先按这个key查数据库如果已存在就直接返回已有任务的task_id和状态不会重复创建。底层就靠数据库唯一索引兜底避免并发场景下的竞态。-- 任务表中加唯一索引 UNIQUE KEY uk_idempotency (user_id, idempotency_key)注意这个唯一索引不能跨用户否则不同用户提交了相同内容的prompt会被误判成重复。我们最开始就犯了这个错索引只建在了idempotency_key上结果两个用户生成一只猫在跑步第二个用户的请求直接被挡住了。2.3 状态表设计除了状态还要留足排查字段任务状态表是整个网关层的数据核心。我直接给出我们线上使用的核心字段不一定全适用但可以当参考CREATE TABLE generation_task ( task_id BIGINT PRIMARY KEY, user_id VARCHAR(64) NOT NULL, idempotency_key VARCHAR(128) NOT NULL, prompt TEXT NOT NULL, status VARCHAR(20) NOT NULL DEFAULT CREATED, priority TINYINT NOT NULL DEFAULT 5, queue_name VARCHAR(32) NOT NULL DEFAULT default, worker_id VARCHAR(64) DEFAULT NULL, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, queued_at DATETIME DEFAULT NULL, started_at DATETIME DEFAULT NULL, finished_at DATETIME DEFAULT NULL, result_url VARCHAR(512) DEFAULT NULL, error_code VARCHAR(32) DEFAULT NULL, error_msg VARCHAR(512) DEFAULT NULL, retry_count INT NOT NULL DEFAULT 0, callback_url VARCHAR(512) DEFAULT NULL, callback_status VARCHAR(20) NOT NULL DEFAULT PENDING, callback_retry INT NOT NULL DEFAULT 0, callback_payload JSON DEFAULT NULL, UNIQUE KEY uk_idempotency (user_id, idempotency_key), KEY idx_status_updated (status, updated_at), KEY idx_callback_status (callback_status, callback_retry) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;这里有三个索引值得说一下。idx_status_updated是为了让调度器快速扫描超时任务比如WHERE status IN (QUEUED,PROCESSING) AND updated_at NOW() - INTERVAL 5 MINUTEidx_callback_status是给回调重试调度器用的专门捞那些回调还没送达的任务。另外queued_at、started_at这些时间戳字段不要偷懒不写后面做耗时分析、性能排查全靠它们。3. 提交链路与队列选型的实际考量3.1 提交接口的处理流程先写库再发消息任务提交接口是整个链路的第一道关卡。我们最终实现的处理顺序是参数校验检查prompt是否为空、长度是否超限、图片是否合法、用户是否还有配额幂等检查根据idempotency_key查库命中就直接返回已有任务创建任务插入一条statusCREATED的记录拿到task_id发送队列消息把task_id和必要参数推送到MQ更新状态如果消息发送成功把status更新为QUEUED如果发送失败把status置为FAILED。这里最核心的一个设计原则是先落库再发消息消息发送成功后才更新状态。为什么要这样大家想一个问题如果先发消息再落库worker那边已经把任务拉起来开始算了数据库里却找不到这个任务状态从何谈起反过来先落库再发消息即使消息发送失败任务至少是CREATED状态调度器可以捞出来重新处理。关于消息发送成功但数据库更新失败的情况我的建议是不要在这个环节做太复杂的事务。我们曾经尝试把插入任务和更新状态放在一个数据库事务里并且把发送MQ也做成事务消息。后来发现事务消息在并发量上来之后要么是发送超时导致整个事务回滚要么是消息发出去了但事务迟迟没提交消费者拿不到数据。很折磨。最终我们简化成了上面这种本地事务更新数据库 消息发出后异步确认的模式配合一个补偿任务每隔几分钟扫描一次CREATED状态超过N分钟的任务重新补发消息。3.2 队列选型Redis List还是MQ文生视频业务对队列的核心要求是容量要大、支持延迟重试、能支撑一定的优先级。我们先后对比了几种方案Redis ListBRPOPLPUSH简单延迟低但不够灵活重试、死信、优先级都不好做而且大任务堆积时Redis内存会告急。RabbitMQ功能完整支持延迟队列、死信队列、优先级队列运维也相对成熟。Kafka吞吐量极高但更适合事件流场景用来做任务队列有点大材小用而且Kafka的按key消费和重试语义不如RabbitMQ直观。我们最终选了RabbitMQ。一个很重要的原因是文生视频的任务量其实没有想象中那么大几十条queue的吞吐完全够用RabbitMQ的重试和TTL机制用起来很舒服。具体上我们建了三个队列task.submit.queue普通任务提交worker从这里消费task.delay.queue延迟重试队列比如worker处理失败后把消息丢到这里3分钟后重新入队task.dlx.queue死信队列处理失败超过N次的消息进死信便于人工捞出来分析。优先级这块RabbitMQ的x-max-priority参数实测有效但别把优先级档位设置太多因为每个优先级其实是独立的内部队列太多优先级反而增加调度开销。我们用3档低1、普通5、高10高优先级一般给付费用户或手动重新生成的任务。3.3 worker消费的坑先改状态还是先处理任务worker消费消息后的顺序竟然也是个翻车点。我们最早是拿到消息就直接开始GPU推理推理完再更新任务状态。看起来没毛病但问题出在如果worker在推理过程中突然宕机消息已经从队列里拿走了数据库里任务还停在QUEUED等worker重启这条消息早就丢了任务就永远卡在QUEUED状态没人去管。这里推荐的做法是worker拿到消息后第一件事就是把数据库状态从QUEUED更新为PROCESSING并且记录worker_id然后再去调GPU推理。这样即使worker中途宕机至少状态是PROCESSING调度器扫描超时任务时就能发现它并重新投递。另外状态更新和消息确认的顺序也要注意。我们总结的经验是先更新数据库状态为PROCESSING再给MQ确认ack。因为一旦ack这条消息就彻底没了如果状态没更新成功任务就处于队列里找不到、数据库还是QUEUED的尴尬阶段只能靠补偿任务捞。先更新状态再ack虽然理论上有更新成功但ack失败导致消息重复消费的风险但重复消费最多是重复生成一次视频我们可以用数据库里的task_id做去重判断比消息丢失更好恢复。这段经验可以总结成一句话异步链路中宁可重复消费也不要消息丢失。重复能靠幂等去重丢失就只能靠扫描补偿而补偿是有时间窗口的用户根本等不了。4. 回调通知与状态查询两条腿走路4.1 回调通知为什么不能做成尽力而为文章开头说的那次事故就是回调丢消息。从那以后我们对回调做了彻底重构。核心思路借鉴了事务发件箱Transactional Outbox模式业务操作和通知事件在同一个本地事务里搞通知事件先落库再通过异步调度器把事件推送给客户端。具体实现是这样的worker完成生成后把result_url写回任务表同时把callback_status置为PENDING表示有一个回调待发送回调调度器每隔几秒扫描callback_statusPENDING的任务发现待发送任务后向客户端的callback_url发一条HTTP POST请求请求体包含task_id、status、result_url、error_info等如果客户端返回2xx把callback_status更新为SUCCEEDED否则记录错误callback_retry加1等待下一轮扫描重试重试达到上限比如5次后把callback_status置为FAILED此时只能靠客户端轮询兜底。这个方案的好处是回调消息不再依赖MQ的可靠性而是躺在自己的数据库表里只要数据库不丢回调就一定能被重试到。和纯MQ方案相比虽然实时性稍差扫描间隔几秒但对文生视频这种用户等几十秒都不在乎的场景完全够用。# 伪代码回调调度器核心逻辑 def process_pending_callbacks(): tasks db.query( SELECT * FROM generation_task WHERE callback_status PENDING AND callback_retry 5 ORDER BY finished_at ASC LIMIT 100 ) for task in tasks: try: resp http.post(task[callback_url], json{ task_id: task[task_id], status: task[status], result_url: task[result_url], error_code: task[error_code], }, timeout10) if resp.status_code 200: db.execute( UPDATE generation_task SET callback_statusSUCCEEDED WHERE task_id?, task[task_id] ) else: db.execute( UPDATE generation_task SET callback_retrycallback_retry1 WHERE task_id?, task[task_id] ) except Exception: db.execute( UPDATE generation_task SET callback_retrycallback_retry1 WHERE task_id?, task[task_id] )注意回调请求的超时不能设太长10秒足够了。如果客户端回调接口响应慢说明对方服务也有问题重试几次就会自动放弃不用死磕。4.2 幂等回调客户端收到两次回调也不怕回调重试带来的副作用是客户端可能会收到重复回调。比如第一次回调其实已经抵达了只是服务端因为网络超时报错于是重试了第二次。如果客户端没有幂等处理看到两条回调就可能会创建两个下载任务、收两次费用。所以回调的请求体里一定要带上task_id并且建议客户端的回调处理接口按task_id做幂等。我们还在响应头里加了Idempotency-Key跟请求体里的task_id保持一致。另外我们约定回调是最终状态通知只有SUCCEEDED或FAILED这两种终态才会触达回调中间状态一律不通知尽量减少客户端的处理复杂度。4.3 查询接口轮询的时机和频率尽管做了可靠的回调通知轮询查询接口依然是必须存在的兜底。因为回调可能因客户端服务故障而最终失败callback_statusFAILED此时用户只能通过轮询拿到结果。状态查询接口做得非常简单GET /api/v1/tasks/{task_id}返回任务当前状态和时间戳。但轮询频率怎么定其实有讲究。我们走过一段弯路前端每500ms轮询一次导致查询接口的QPS是提交接口的几十倍给数据库造成了不小的压力。后来我们调整策略让前端做自适应轮询任务刚提交时每2秒查一次超过30秒没结果时改成每5秒查一次超过2分钟后改成每10秒查一次。从网关层角度我们还在查询接口上做了short-circuit优化如果任务已经是终态SUCCEEDED或FAILED查询结果直接走本地缓存不再查数据库因为终态数据几乎不会再变化。这个优化把查询接口的数据库压力降了大概70%。5. 超时、限流与任务取消的治理细节5.1 超时设计排队超时和处理超时分开文生视频任务超时是很容易被忽略但又很致命的问题。我们定义了两种超时排队超时任务进入QUEUED后如果在N分钟内没有被任何worker消费视为排队超时。我们的排队超时设为3分钟因为高峰期确实会有排队这个值给得比较宽。处理超时任务进入PROCESSING后如果在M分钟内没有变成终态视为处理超时。视频生成的耗时跟prompt长度、分辨率、帧数都有关系我们最初设为3分钟后来发现高清视频偶尔要4分钟所以调到5分钟。超时检测靠一个定时调度器完成每分钟扫描一次SELECT task_id FROM generation_task WHERE status QUEUED AND queued_at IS NOT NULL AND updated_at NOW() - INTERVAL 3 MINUTE; SELECT task_id FROM generation_task WHERE status PROCESSING AND started_at IS NOT NULL AND updated_at NOW() - INTERVAL 5 MINUTE;扫描到超时任务后不能直接置成FAILED就完事。我们的处理是排队超时的任务先尝试重新投递到队列重试2次后还是超时才置为FAILED错误码写成QUEUE_TIMEOUT。处理超时的任务先检查worker节点是否还活着通过worker心跳表如果worker活着可能是生成慢再多给一点宽限时间如果worker已经失联立即把任务重新入队。5.2 限流从入口控制并发文生视频的GPU资源是硬约束Active状态的任务数量一旦超过GPU卡数剩下所有任务都会积压。所以网关层的限流要同时cover几个维度用户维度限流每个用户同一时刻最多N个未完成任务。我们设的是3超出直接返回TOO_MANY_PENDING_TASKS让用户等已有任务完成再提交。全局并发限流全系统同时处于PROCESSING状态的任务数不能超过GPU卡数乘以一个系数比如1.5留一点buffer。超出的任务卡在QUEUED阶段不进入worker。这样即使有人恶意刷接口GPU也不会被打爆。提交速率限流每用户每分钟最多提交M次防脚本刷。限流发生在提交接口的最前面先限流再走幂等和落库否则恶意请求会把数据库打垮。我们用的就是一个简单的Redis计数器INCR加过期时间代码几行就搞定效果很稳定。5.3 用户取消终态之前怎么处理任务取消也是个容易出问题的点。用户点了取消如果任务还在排队中QUEUED好办直接更新状态为CANCELED再从队列里把消息捞出来丢掉即可。但如果在PROCESSING状态GPU已经在算视频了这时候取消有两个选择硬取消直接告诉worker任务取消了让它停止推理释放GPU。省资源但有可能生成一半的视频文件残留需要清理。软取消让worker继续算完但结果不再通知用户只在后台保留。浪费资源但逻辑简单。我们最终的策略是默认软取消用户可以主动选立即停止来走硬取消。因为从产品角度用户取消后往往还会回来再生成一次软取消让worker算完其实下一次的生成结果可以直接复用。不过这只是我们的产品取舍不一定适合所有场景大家根据自己的成本模型来定。6. 常见问题与排查技巧实录6.1 高频故障速查表我把这段时间遇到的高频问题整理成了一个排查速查表遇到同类问题可以直接对号入座故障现象可能原因排查顺序解决方案任务一直QUEUED没人消费worker挂了1. 查worker心跳; 2. 查MQ消费者列表; 3. 查消息是否堆积重启worker重新投递消息任务一直PROCESSINGworker进程僵死1. 查worker日志; 2. 查GPU显存使用率kill掉僵死进程调度器重新入队回调一直PENDING不触发回调调度器停了1. 查调度器日志; 2. 查扫描SQL是否有阻塞重启调度器检查SQL索引客户端收到重复回调回调重试机制1. 查回调日志的timeout; 2. 查客户端幂等逻辑客户端按task_id去重用户反馈提交没反应幂等键丢失或冲突1. 查前端是否传幂等键; 2. 查数据库唯一索引冲突前端统一生成幂等键高峰期提交延迟高限流触发或数据库慢查询1. 查Redis限流计数; 2. 查DB慢SQL扩大限流阈值优化索引6.2 排查状态不一致的三个技巧异步任务最怕的是数据库状态和实际业务状态对不上。排查这种问题我总结出三个比较实用的技巧。第一看时间戳。排到状态卡住的任务时先看queued_at、started_at、finished_at这几个时间戳一对比很快能判断卡在哪个环节。比如started_at不为空但状态还是PROCESSING且超过5分钟那基本就是worker处理中异常退出。第二看worker_id。PROCESSING状态一定要记录worker_id排查时拿着worker_id去查这台机器的日志和心跳能确认这台机器是否还活着。没有worker_id的话这个问题会变得非常难查只能盲猜。第三回放补偿。不要一次性把状态不对的任务全部重置稳妥的做法是先查出来按任务ID小批量处理比如一次50条处理完跑一遍状态统计确认没问题再处理下一批。大范围直接update是一次高风险操作很容易把还正常跑着的任务也带偏。6.3 数据补偿脚本的教学示例最后给一个非常实用的补偿脚本思路。当调度器因为各种原因漏处理了一批卡住的任务时你需要一个人工或半自动的补偿流程。我们用的就是一个Python脚本配合SQL查询先把候选任务捞出来再逐个决策#!/usr/bin/env python3 # 补偿脚本扫描超时未完成的任务并重置 import mysql.connector conn mysql.connector.connect(...) def find_stuck_tasks(limit100): cur conn.cursor(dictionaryTrue) cur.execute( SELECT task_id, status, created_at, started_at, worker_id FROM generation_task WHERE status IN (QUEUED, PROCESSING) AND updated_at NOW() - INTERVAL 10 MINUTE ORDER BY created_at ASC LIMIT %s , (limit,)) return cur.fetchall() def requeue_task(task_id): # 重新投递到MQ这条消息需要包含完整task信息 mq.publish(task.submit.queue, {task_id: task_id}) for t in find_stuck_tasks(50): if t[status] PROCESSING: # 先确认worker是否存活如果存活且生成中跳过 if worker_alive(t[worker_id]): continue requeue_task(t[task_id]) print(frequeued task {t[task_id]})注意补偿脚本一定要有--dry-run参数先跑一遍只打印不执行确认无误后再真正执行。我刚开始写补偿脚本的时候也在生产上吃过亏一激动直接跑了个update把正常任务的状态都改乱了后来老老实实先dry-run。7. 写在最后一点实战体会回看这一路文生视频业务的异步任务网关层最难的不是技术选型而是把任务全生命周期这个思维贯穿到每一个细节里。状态机、幂等、回调、超时、限流单独拎出来哪一个都是老话题但真正合在一起面对真实业务的时候总会有不少意料之外的暗坑。最大的体会就一句话异步任务网关层的核心不是快而是可靠是让每一个任务都有一条清晰可见、可追踪、可补偿的生命周期路径。如果你正在做类似的AIGC视频业务我的建议是前期宁可多写几张表和几个调度器也不要在可靠性和可观测性上省钱。那些看起来麻烦的机制比如outbox回调、超时扫描、补偿任务几乎都是等线上出过事故之后才补上的早期一旦缺失后期补的成本会高很多。最后再分享一个小技巧每个关键状态流转的地方都打一条结构化日志带上task_id和状态排查问题的时候一条链路从头拉到尾比任何监控面板都好用。
返回列表