Post

恢复任务,为什么会拖慢新任务?

恢复旧任务为什么会挤占新任务的执行机会?从异步等待、批次处理与资源竞争入手,分析恢复链路影响正常吞吐的原因。

项目 阅读 18 点赞 0 评论 0

异步系统里的等待、批次与资源竞争

一个后台系统可能没有丢任务,却仍然让用户等得很不合理。

请求已经写入数据库,队列也没有堆积,worker 看起来并不忙,用户刚提交的任务却迟迟没有开始。此时继续增加 worker,可能几乎没有效果:消息还没有交到它们手里。

ProductFlow 的后台要处理商品图片生成、图片交付和其它异步工作。除了发送新任务,dispatcher 还会检查历史上未完成的任务,判断哪些需要补投、哪些仍在运行、哪些应结束为结果不明。恢复本来是为了帮助系统继续前进,但如果安排不当,它本身也会制造等待。

这篇文章讨论一个容易被“增加并发”掩盖的问题:系统里不同来源的等待,应该分别在哪里解除? 程序调用顺序、轮询周期、批次边界和数据库锁,都可能让任务停住;它们需要的修改并不相同。

1 先弄清楚用户等在哪一段

一次生图的总耗时包含很多阶段:

接受请求
  → 数据库中的投递等待
  → 队列等待
  → Worker 执行准备
  → 图片供应商生成
  → 结果保存与展示

本文关注第二段:待投递记录已经提交,到 dispatcher 将消息交给队列。这个范围看似很窄,却决定了“任务为什么还没开始”是否能够得到正确解释。

如果把所有时间揉成一个“生成耗时”,供应商的长尾会掩盖调度问题,调度的周期性空等也可能被误归因到模型速度。平均队列长度同样不够:消息尚未进入队列时,队列越空也不能证明系统越快。

一个最小的观测,需要至少区分任务提交、投递成功与 worker 开始三个时刻。时钟与统计口径还要保持一致,避免把不同机器的时间误差当作几毫秒级的优化收益。

2 一行调用顺序,就能形成跨业务等待

2.1 串行循环的隐含依赖

最直观的 dispatcher 可以写成这样:

每轮:
    恢复历史任务
    投递新的到期任务
    等待下一轮

这个结构容易理解,也能在任务少、恢复快时正常工作。但它让新任务依赖整轮恢复完成。

设想恢复正在检查一批旧图片任务,其中一条查询长时间等待数据库。另一个商家操作提交了图片导出,业务已经接受,投递代码却还没执行。两项工作没有业务依赖,只因为函数排在后面,导出就多等了一段时间。

在一个简化模型中,投递前的等待可以分成:

等待下一轮开始
+ 排在投递前面的恢复时间
+ 查询与资源竞争时间

缩短轮询周期能影响第一项,增加 worker 影响后续执行能力,都不能移除第二项。定位到这个因果点之后,修改才会明确:让正常投递拥有独立的调度循环。

2.2 错误隔离与等待隔离是两回事

即使恢复函数已经做到“某个域报错后继续处理其它域”,仍然可能被一个迟迟不返回的调用挡住。

错误处理发生在函数返回之后。一个慢函数尚未返回时,后续域根本没有机会开始。这也是为什么“错误被捕获了”不能证明“其它业务不会受影响”。

ProductFlow 后来将常驻模式下的正常投递与各业务域恢复分别调度。当前涉及工作流、连续生图、交付、局部编辑和 Agent,各域内部串行处理本批候选。

常驻 Dispatcher

正常投递循环

工作流恢复

连续生图恢复

交付恢复

局部编辑恢复

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 锁定子句文档

正常 Worker数据库恢复器正常 Worker数据库恢复器短事务发现候选返回 ID 并释放锁完成其中一个任务逐项重查当前状态已完成,无需恢复

图 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 scanInnoDB 索引结构

队列状态变化频繁,索引也一直在维护。给更多列建索引会增加写入和空间成本。这条链路需要的是对准领取与恢复查询的少数索引,而不是把所有可查询字段都加上一遍。

7 拆开循环以后,瓶颈可能转移到连接池

程序依赖解除后,几个域能够同时查询,也就可能同时争用有限连接。新增并发会提高可调度机会,实际吞吐还取决于事务时长、连接池、锁和存储性能。

因此,观察到延迟下降以后,仍应同时看任务年龄与资源成本。一个指标很好看,可能只是压力被移到了另一处。

同时出现的现象 更可能需要调查的地方
数据库资源空闲,延迟呈周期性台阶 轮询唤醒、批次之间是否仍等待
恢复调用未返回,新投递也不开始 是否仍有串行依赖或共享互斥
各域同时变慢,连接等待增加 事务占用与连接预算
每批耗时很短,最老任务持续变老 跳锁、候选选择与反复失败
Redis 队列很空,待投递记录却很多 Dispatcher 的发送能力和错误
新任务变快,旧任务几乎不再推进 恢复投入不足或公平性问题

进一步独立连接池、错峰调度或拆进程,都应该对应这里发现的瓶颈。没有连接等待证据就增加隔离,只会增加资源与配置责任。

对于容量规划,可以用一个粗略关系思考:某域每轮最多处理 B 条,平均一轮实际耗时为 T,那么长期处理率受 B/T 限制;如果恢复候选产生得更快,单靠等待更久不会消化积压。实际 T 包含处理与调度,并不能简单拿配置 interval 代替。

8 用确定性时序测隔离,用负载测收益

调度关系很适合用可控回调验证:让某个恢复域停在 channel 上,保持它没有返回,同时检查通知是否触发了投递、其它恢复域是否开始工作。然后发出取消,释放测试回调,观察循环退出。

这种测试比“睡两秒再检查结果”更容易暴露真实依赖,也不会把机器偶然快慢当作正确性。

当前调度测试覆盖了恢复阻塞时投递仍可被通知或 ticker 唤醒、跨域继续执行、积压续批,以及取消后不启动新批次等行为。

数据库锁需要真实事务验证。可以锁住排序最前的一批候选,让后面保留一条可处理任务;恢复应能推进后者,释放锁以后再检查原先的任务。只用 mock 返回候选数组,无法证明真实 SQL 避开了队首阻塞。

性能收益则需要另一组实验:固定任务规模、连接池、并发到达率和故障注入方式,对照投递延迟分布、最老候选年龄、每域有效恢复数与数据库等待。本文没有报告这样的同条件性能数据,不能把消除某条程序依赖写成已证明的吞吐倍数。

恢复系统尤其需要同时测“新任务还快不快”和“旧任务最终有没有推进”。只追求其中一项,可能把另一项饿住。

9 每种等待都要有自己的解释

这套调整没有一个统一的“加速开关”。独立循环解除函数顺序造成的等待,通知减少轮询空等,满批续投减少批次间隔,跳锁让当前可处理的候选有机会进入本轮。数据库资源竞争与长期公平性仍然存在。

对用户而言,这些差别应该表现为更直接的结果:刚提交的工作不必陪历史扫描一起等待,旧任务也不会因为新请求不断到来而永远被遗忘。

对工程师而言,最有用的习惯是沿时间线追问:任务现在停在哪里,前面在等谁,解除这段等待后又会占用什么资源。这个问题比“还能不能多开几个协程”更容易找到有效的修改。

延伸阅读

继续阅读

继续阅读

全部归档

评论