异步系统里的等待、批次与资源竞争
一个后台系统可能没有丢任务,却仍然让用户等得很不合理。
请求已经写入数据库,队列也没有堆积,worker 看起来并不忙,用户刚提交的任务却迟迟没有开始。此时继续增加 worker,可能几乎没有效果:消息还没有交到它们手里。
ProductFlow 的后台要处理商品图片生成、图片交付和其它异步工作。除了发送新任务,dispatcher 还会检查历史上未完成的任务,判断哪些需要补投、哪些仍在运行、哪些应结束为结果不明。恢复本来是为了帮助系统继续前进,但如果安排不当,它本身也会制造等待。
这篇文章讨论一个容易被“增加并发”掩盖的问题:系统里不同来源的等待,应该分别在哪里解除? 程序调用顺序、轮询周期、批次边界和数据库锁,都可能让任务停住;它们需要的修改并不相同。
1 先弄清楚用户等在哪一段
一次生图的总耗时包含很多阶段:
接受请求
→ 数据库中的投递等待
→ 队列等待
→ Worker 执行准备
→ 图片供应商生成
→ 结果保存与展示
本文关注第二段:待投递记录已经提交,到 dispatcher 将消息交给队列。这个范围看似很窄,却决定了“任务为什么还没开始”是否能够得到正确解释。
如果把所有时间揉成一个“生成耗时”,供应商的长尾会掩盖调度问题,调度的周期性空等也可能被误归因到模型速度。平均队列长度同样不够:消息尚未进入队列时,队列越空也不能证明系统越快。
一个最小的观测,需要至少区分任务提交、投递成功与 worker 开始三个时刻。时钟与统计口径还要保持一致,避免把不同机器的时间误差当作几毫秒级的优化收益。
2 一行调用顺序,就能形成跨业务等待
2.1 串行循环的隐含依赖
最直观的 dispatcher 可以写成这样:
每轮:
恢复历史任务
投递新的到期任务
等待下一轮
这个结构容易理解,也能在任务少、恢复快时正常工作。但它让新任务依赖整轮恢复完成。
设想恢复正在检查一批旧图片任务,其中一条查询长时间等待数据库。另一个商家操作提交了图片导出,业务已经接受,投递代码却还没执行。两项工作没有业务依赖,只因为函数排在后面,导出就多等了一段时间。
在一个简化模型中,投递前的等待可以分成:
等待下一轮开始
+ 排在投递前面的恢复时间
+ 查询与资源竞争时间
缩短轮询周期能影响第一项,增加 worker 影响后续执行能力,都不能移除第二项。定位到这个因果点之后,修改才会明确:让正常投递拥有独立的调度循环。
2.2 错误隔离与等待隔离是两回事
即使恢复函数已经做到“某个域报错后继续处理其它域”,仍然可能被一个迟迟不返回的调用挡住。
错误处理发生在函数返回之后。一个慢函数尚未返回时,后续域根本没有机会开始。这也是为什么“错误被捕获了”不能证明“其它业务不会受影响”。
ProductFlow 后来将常驻模式下的正常投递与各业务域恢复分别调度。当前涉及工作流、连续生图、交付、局部编辑和 Agent,各域内部串行处理本批候选。
图 1. 各域可以独立推进,数据库资源仍由它们共同使用。
一次性命令保留恢复后投递的顺序,并汇总错误;常驻服务则需要持续接收新工作,因此拆成独立循环。
3 为什么没有给每个恢复任务都开一个 goroutine
按域拆循环以后,很自然会继续想到:既然并发有效,为什么不把所有历史任务一起恢复?
因为候选发现速度与实际处理能力不同。把一千个任务立即放进 goroutine,只是让它们更快地进入连接池等待队列;数据库能够同时完成多少事务,并没有因此增加。
恢复期间,系统往往已经积累了异常和积压。如果发现多少就启动多少,恢复器会在资源最紧张时制造更大的需求峰值。
按域建立有界串行循环,是当前采用的折中:
| 并发粒度 | 获得的前进机会 | 仍然存在的等待或成本 |
|---|---|---|
| 所有恢复串行 | 执行顺序简单 | 一个慢域拖住后续全部域 |
| 按域独立,域内串行 | 不同业务可分别推进 | 同域慢任务仍影响本域后续候选 |
| 每个候选立即并发 | 候选可以同时发起处理 | 连接、锁、内存与调度竞争迅速增加 |
| 按资源预算限制并发 | 可以针对瓶颈控制投入 | 需要明确容量、成本和公平性模型 |
我们采用第二种,因为它直接对应已有业务恢复入口,可以解除已知的跨域等待。第四种值得在资源竞争成为主要问题时继续研究,但需要知道限制的是数据库连接、外部调用还是本地计算。
独立循环本身也不能保证及时退出。如果某个回调忽略取消并永久阻塞,关闭时等待它结束仍可能挂住。调度器检查取消,只能保证自己不继续启动新工作,已经开始的 I/O 还需要相应超时与取消支持。
3.1 并发限制放在 go 语句的哪一边
正常投递与恢复还有一个不同点:一批已经领取的信封可以并发发送。这里使用容量有限的 channel 作为信号量,但名额是在启动 goroutine 之前取得的。
两种写法表面上都限制了同时工作的数量:
// 写法 A:先启动,再排队。
go func() {
slots <- struct{}{}
defer func() { <-slots }()
send()
}()
// 写法 B:先取得名额,再启动。
slots <- struct{}{}
go func() {
defer func() { <-slots }()
send()
}()
差别在等待发生的位置。A 可以很快创建大量 goroutine,每个都保留自己的栈和引用对象,然后一起等名额;B 让提交方在容量耗尽时停下来,工作数量受到约束。当前投递采用 B,并用 WaitGroup 等待本批发送结束。
goroutine 能降低等待的开销,因为可轮询的网络 I/O 挂起时,运行时可以让线程继续执行其他 goroutine。但数据库连接、Redis 处理能力和内存不会随 go 语句一起增加。只有把压力传回提交方,才能避免“并发有限,等待者无限”。
发送成功数则由 atomic.Int64 累加。这里保护的是一个独立计数;普通 sent++ 包含读取、加一和写回,两个 goroutine 可能读到同一个旧值,最后只留下一个增量。原子加法把这一步变成不可被另一加法拆开的操作。它适合计数,业务任务的多个字段仍然交给数据库事务一起更新。
4 通知只负责唤醒,任务仍然留在数据库
4.1 固定轮询会给每个新请求增加空等
即使没有恢复阻塞,新任务也可能刚好错过一次扫描,只能等待下一轮。
假设扫描间隔为 I,任务在间隔内均匀到达,忽略扫描耗时和竞争,平均轮询等待约为 I/2。即使系统完全空闲,这一段等待仍然存在。
缩短 I 可以改善延迟,却会增加空扫描。ProductFlow 使用 PostgreSQL 通知,在持久投递意图产生后唤醒 dispatcher,同时保留定时扫描。
按照 PostgreSQL 的 NOTIFY 语义,事务内发出的通知只会在事务提交后交付;回滚的事务不会发布这份通知。这使通知与业务提交具有明确顺序。
4.2 为什么容量为一的通道可以工作
当前实现将通知合并到容量为一的唤醒通道。如果已有一个唤醒信号未处理,后来的信号无需继续堆积。
这里能合并,是因为信号的含义是“值得查一次数据库”。任务身份和待执行意图已经持久化,下一次扫描可以找到多项工作。
若每条信号携带唯一任务,而且系统没有其它持久记录,同样的合并就会丢工作。容量为一没有普遍正确性,它依赖信号与任务记录之间的分工。
监听不可用时,定时扫描仍提供发现路径。通知改善通常情况下的等待,但不能成为请求能够被处理的唯一前提。
通知还解决不了一件事:一次扫描可能只拿到一批。剩下的任务如果继续等待下一次 tick,批次之间仍然会不断重复空等。
5 投递满批继续,恢复满批却未必应该继续
5.1 批次限制可能意外变成吞吐上限
设有 N 条到期信封,每批最多处理 B 条,没有新任务进入,也没有错误和资源竞争。需要处理的批次数为:
K = ceil(N / B)
如果每批之后都等待 I,批次之间会增加约 (K - 1) × I 的人为间隔。这里没有计算查询与发送耗时,公式只标出一项由调度结构产生的额外等待。
当前投递在 HasMore 为真时继续取下一批,直到没有更多工作、发生错误或收到取消;单批发送并发也有限制。
它移除了已知积压之间的空等。若持续到达率超过处理能力,积压仍会增长;续批无法解决容量不足。
5.2 恢复的收益与投递不同
投递面对的是已经到期、需要交给执行器的工作。恢复面对的是需要重新检查的历史状态;一次扫描未必产生任何可执行动作。
例如某项外部调用仍在等待,或者候选发现后任务已经被 worker 推进,检查结束也可能只是跳过。如果恢复满批后无限续扫,就可能在故障期间持续投入大量数据库资源,却没有相称的有效恢复。
因此当前恢复保留按域周期调度与有界批次,没有复用投递的满批续取循环。
图 2. 同样按批处理,调度策略可以不同。 恢复不根据满批结果主动连续扫完整个积压;实际间隔还受批次耗时和 ticker 行为影响。
这个折中接受了更长的历史积压处理时间,换取较可控的恢复投入。如果某个恢复域不断积累老任务,需要查看有效恢复率和候选年龄,再决定调整批次、节奏或并发。
6 跳过锁住的行,后面的工作才有机会
循环已经独立,唤醒也已经及时,数据库查询仍可能把它挡住。
假设队首的一批旧任务被长事务锁定。普通加锁读取会等待,后面原本可以处理的候选也进不了本轮结果。对于队列式候选发现,PostgreSQL 的 SKIP LOCKED 允许跳过当前无法取得行锁的记录。参见PostgreSQL 锁定子句文档。
ProductFlow 的交付恢复采用短事务发现候选,按更新时间和 ID 排序,最多读取批次额度加一条。当前默认额度为 25,多读的一条用于判断本次可见候选是否超过额度。
-- 说明查询结构的简化示例
BEGIN;
SELECT id
FROM image_jobs
WHERE status = 'queued'
ORDER BY updated_at, id
LIMIT 26
FOR UPDATE SKIP LOCKED;
COMMIT;
这里有三个必须一起理解的限制。
候选身份不等于处理权。 发现事务结束后,锁已经释放。真正恢复时要重新读取状态:正常 worker 可能已经完成它,用户也可能取消了它。候选列表只是值得继续检查的对象。
HasMore 不等于积压总量。 它只描述本次跳过锁定行以后看见的候选。所有老任务都被锁住时,返回很少的记录并不能证明系统没有积压。
跳锁不保证公平。 一条反复被占用的记录可能长期被跳过。SQL 返回得很快,与最老任务终于得到处理,是两件不同的事。并且 SKIP LOCKED 只改变相应行锁的等待,不能解除全部数据库阻塞。参见PostgreSQL 锁定子句文档。
图 3. 发现与处理之间允许正常工作继续。 重查将时间差纳入协议,避免把陈旧候选当成必须执行的命令。
6.1 跳过锁之前,先高效地找到该锁哪几行
投递查询按状态筛选,再按到期时间和 ID 排序。省略租约条件后,主要形状如下:
SELECT id
FROM async_dispatches
WHERE status = 'PENDING' AND available_at <= :now
ORDER BY available_at, id
LIMIT 100;
对应的联合索引是 (status, available_at, id)。它相当于一份先按状态分区、再按时间和编号排列的目录:查询先找到 PENDING 区间,沿时间顺序取出到期条目,同一时刻再用 ID 决定顺序。
为什么不分别给三个字段建索引?单列索引可以帮助各自的筛选,却未必同时提供我们需要的组合顺序。数据库可能合并候选以后继续排序,也可能认为直接扫描更划算。联合索引的价值,是让一次访问尽量顺着这条查询的工作方式进行。
这里可以看到 B-tree 索引适合数据库的原因。数据库访问以页为单位,一个内部页可以容纳很多导航条目,树的分支多,就能用较少的层级定位到叶层;到达叶层后,还能沿排序范围继续读取。平衡二叉树虽然也是对数查找,每次指针跳转提供的分支却少得多。哈希表善于等值查找,无法直接给出“截止现在最早的 100 项”。
字段顺序同样来自查询。如果时间放在状态前面,所有状态的旧任务会先按时间混在一起,查询可能需要越过更多无关记录。范围条件以后的列也不是突然毫无用处:它仍可能参与过滤或满足排序,只是未必继续缩小同一个连续扫描区间。具体要看数据库版本和联合索引规则。
索引解决的是定位成本。领取任务还要访问行、核对租约并取得锁,SKIP LOCKED 解决的才是其中的行锁等待。因此,优化查询时要同时看扫描量、过滤量、排序和锁等待,只看“用了索引”很容易漏掉真正耗时的地方。
6.2 回到表里读取的成本
ProductFlow 使用 PostgreSQL,普通表数据保存在堆中,索引条目定位相应元组。即便所需字段都在索引里,是否能省掉堆访问,还要看可见性信息;领取任务需要加行锁,也不能只在索引中完成。
这与 InnoDB 常见的解释有所区别:InnoDB 的聚簇索引保存行数据,二级索引带着主键,必要时再按主键定位行。把两者都简称为“回表”很方便,但分析一次具体查询时,仍需要知道到底多访问了什么。可分别参照 PostgreSQL 的 index-only scan与 InnoDB 索引结构。
队列状态变化频繁,索引也一直在维护。给更多列建索引会增加写入和空间成本。这条链路需要的是对准领取与恢复查询的少数索引,而不是把所有可查询字段都加上一遍。
7 拆开循环以后,瓶颈可能转移到连接池
程序依赖解除后,几个域能够同时查询,也就可能同时争用有限连接。新增并发会提高可调度机会,实际吞吐还取决于事务时长、连接池、锁和存储性能。
因此,观察到延迟下降以后,仍应同时看任务年龄与资源成本。一个指标很好看,可能只是压力被移到了另一处。
| 同时出现的现象 | 更可能需要调查的地方 |
|---|---|
| 数据库资源空闲,延迟呈周期性台阶 | 轮询唤醒、批次之间是否仍等待 |
| 恢复调用未返回,新投递也不开始 | 是否仍有串行依赖或共享互斥 |
| 各域同时变慢,连接等待增加 | 事务占用与连接预算 |
| 每批耗时很短,最老任务持续变老 | 跳锁、候选选择与反复失败 |
| Redis 队列很空,待投递记录却很多 | Dispatcher 的发送能力和错误 |
| 新任务变快,旧任务几乎不再推进 | 恢复投入不足或公平性问题 |
进一步独立连接池、错峰调度或拆进程,都应该对应这里发现的瓶颈。没有连接等待证据就增加隔离,只会增加资源与配置责任。
对于容量规划,可以用一个粗略关系思考:某域每轮最多处理 B 条,平均一轮实际耗时为 T,那么长期处理率受 B/T 限制;如果恢复候选产生得更快,单靠等待更久不会消化积压。实际 T 包含处理与调度,并不能简单拿配置 interval 代替。
8 用确定性时序测隔离,用负载测收益
调度关系很适合用可控回调验证:让某个恢复域停在 channel 上,保持它没有返回,同时检查通知是否触发了投递、其它恢复域是否开始工作。然后发出取消,释放测试回调,观察循环退出。
这种测试比“睡两秒再检查结果”更容易暴露真实依赖,也不会把机器偶然快慢当作正确性。
当前调度测试覆盖了恢复阻塞时投递仍可被通知或 ticker 唤醒、跨域继续执行、积压续批,以及取消后不启动新批次等行为。
数据库锁需要真实事务验证。可以锁住排序最前的一批候选,让后面保留一条可处理任务;恢复应能推进后者,释放锁以后再检查原先的任务。只用 mock 返回候选数组,无法证明真实 SQL 避开了队首阻塞。
性能收益则需要另一组实验:固定任务规模、连接池、并发到达率和故障注入方式,对照投递延迟分布、最老候选年龄、每域有效恢复数与数据库等待。本文没有报告这样的同条件性能数据,不能把消除某条程序依赖写成已证明的吞吐倍数。
恢复系统尤其需要同时测“新任务还快不快”和“旧任务最终有没有推进”。只追求其中一项,可能把另一项饿住。
9 每种等待都要有自己的解释
这套调整没有一个统一的“加速开关”。独立循环解除函数顺序造成的等待,通知减少轮询空等,满批续投减少批次间隔,跳锁让当前可处理的候选有机会进入本轮。数据库资源竞争与长期公平性仍然存在。
对用户而言,这些差别应该表现为更直接的结果:刚提交的工作不必陪历史扫描一起等待,旧任务也不会因为新请求不断到来而永远被遗忘。
对工程师而言,最有用的习惯是沿时间线追问:任务现在停在哪里,前面在等谁,解除这段等待后又会占用什么资源。这个问题比“还能不能多开几个协程”更容易找到有效的修改。
评论