几个月前我接了一个内部系统的改造任务,需求看起来非常简单:上游每天会推过来几十万条带分组标识的任务数据,下游接口对并发有硬性限制——同一个分组内的请求必须串行处理,不同分组之间可以并行。我第一版直接用线程池硬写,结果跑了一周就发现问题不断:某个分组处理得慢,整个池子跟着被拖垮;下游限流一触发,所有请求一起遭殃。于是我用几天时间写了一个小工具,取名 Bgrp,全称 Batch Group Processing,也就是"批量分组处理"。这个实验本身不大,但把并发调度里最容易踩的坑几乎都踩了一遍。借着这篇文章,我把设计思路、核心代码、实测数据和踩坑记录完整整理出来,给正在写任务调度、消息消费、批量推送这类逻辑的同学做个参考。
1. 为什么会有这个实验:被"分组并发"折磨过的人都懂
1.1 需求场景:看似简单,实则处处是约束
我当时要处理的是一个商户运营通知的推送系统。上游每天凌晨会把当天需要推送的通知数据打到一个中间表里,我的服务要扫出来,按照商户维度调下游的推送接口。数据量大概每天三十万到五十万条,涉及两百多个商户。约束条件表面上只有两条:同一个商户的推送必须按时间先后顺序执行;不同商户之间可以同时推送。但往下细挖,真实约束远比这两条多。
第一,下游接口的单商户并发上限是 1,也就是说同一个商户同时只能有一个请求在途,多发了就报错。第二,下游整体还有一个 TPS 上限,虽然不像单商户那么严,但并发超过一定量就会触发限流,而且限流之后是直接拒绝,不会排队。第三,任务本身允许重试,但重试不能破坏顺序,否则后发的通知先到、先发的通知反而没到,用户收到的通知顺序就乱了。
这三条叠加在一起,就从一个简单的"循环调用"问题,变成了一个带顺序约束、并发约束和背压约束的调度问题。我最初觉得这不就是个线程池加队列嘛,后来发现完全不是一回事。
1.2 第一版线程池方案的三个问题
第一版我用了ThreadPoolExecutor(20)加上一个全局Queue,提交任务的时候丢进队列,工作线程从队列里取任务直接调用下游接口。表面看逻辑没问题,压测一跑,三个问题立刻暴露出来。
第一个问题是队头阻塞。某个商户的下游接口因为参数问题每次都超时,超时重试要占住一个工作线程好几秒,二十个线程很快就被几个"慢商户"占满了,其他商户的通知全部排队干等。这是典型的 Head-of-Line Blocking,单点故障通过共享线程池被放大成了整体故障。
第二个问题是组内乱序。同一个商户的两条通知如果被两个工作线程同时取走,先取到的线程可能因为下游响应慢,反而比后取到的线程晚完成。数据层面就会出现通知 B 先发出、通知 A 后发出的情况。我那时候做的"顺序保证"沦为空谈。
第三个问题是背压不可控。全局队列用无界队列,任务只会往里堆,下游一旦挂掉,内存里的积压任务越堆越多。如果换成有界队列,满了之后是拒绝还是阻塞,又得自己设计一套策略,代码越写越复杂,行为还未必符合预期。
1.3 这个实验适合谁参考
如果你正在写下面这类逻辑,Bgrp 这个实验应该能帮到你:消息队列的消费者,尤其是 Kafka 这种分区有序、跨分区并行的消费场景;批量推送、批量通知、批量对账这类按业务维度拆分的任务系统;需要调用第三方接口、对单账户或单商户有并发限制的对接层;以及任何"既要组内串行、又要组间并行、还要控制整体并发"的调度需求。
实验的产出是一个两三百行的 Python 小框架,核心思路不绑定语言,换成 Go、Java 也完全可以复刻。后面几节我会把模型的拆解、代码的关键细节、以及跑出来的实测数据都讲清楚。
2. Bgrp 的核心模型:把并发问题拆成两个独立维度
2.1 组是业务边界,批是执行单元
很多人写并发调度,上来就盯着"线程"和"队列",这其实是把问题的层次搞错了。Bgrp 的核心想法是先把约束拆开:组(Group)决定谁不能同时跑,批(Batch)决定怎么跑才高效。
组是业务概念的映射。同一个商户是一条业务链路,链路内部必须串行;同一个 Kafka 分区是一条有序流,分区内部必须串行;同一个仓库的库存变更是一组强一致操作,也必须串行。这些"必须串行"的根源来自业务,不是来自线程。所以在 Bgrp 里,组是一个逻辑隔离单元,每个组有一个独立的缓冲区,同一时刻最多只有一个批次在途。
批是执行层面的产物。任务到达之后不直接调度,而是先在组缓冲区里攒着,调度器按条件把它们打包成一个个批次,再丢给线程池执行。批的大小直接决定了吞吐的底子:单条调用一次网络请求,和一百条调用一次批量接口,开销差一个数量级;即便下游不支持批量接口,把一百个任务打包后连续发出,也能减少线程切换和调度唤醒的次数。
把这两个维度拆开之后,原来那个"线程池加全局队列"的方案就变成了两层调度:第一层是组级别的排队,决定哪些组可以出队;第二层是批级别的并发,决定同时有几个批次在跑。
2.2 两层调度的数据流与组内串行实现
Bgrp 的数据流可以概括成四步:提交、攒批、派发、执行。
- 提交:外部调用
submit(group_key, task),任务进入对应组的 deque 缓冲区。 - 攒批:调度线程每隔一个固定时间醒来,扫描所有非空且不在执行中的组。
- 派发:从选中的组缓冲区里一次性弹出最多
batch_size个任务,组成一个批次,放进线程池。 - 执行:批次在工作线程里逐个执行任务,执行完毕后把该组标记为"空闲",允许下一个批次被派发。
组内串行的关键就是那个"不在执行中"的标记。一个组一旦有批次在跑,调度器就不会再从它里面取任务;只有整个批次跑完,标记清除,后续任务才能继续派发。这样一来,同一组的任务天然被串行化,不需要额外加锁,也不存在多线程同时操作同一组数据的问题。
我拿这个模型和我最初那个"全局队列 + 线程池"的方案对比过,本质区别在于:全局队列方案把"任务应该被谁执行"作为唯一维度,而 Bgrp 把"哪些任务不能同时执行"这个约束显式建模成了"组"这个概念。约束永远应该是逻辑层面的东西,不该靠运行时的偶然行为去凑。
2.3 为什么不用现成的 ThreadPoolExecutor 加队列直接改
这是一个我经常被问到的问题:既然ThreadPoolExecutor和Queue都是现成的,为什么不直接在上面加几行代码?
原因很简单:ThreadPoolExecutor只解决"并发执行"的问题,不解决"哪些任务互斥"的问题。你要在它上面实现组内串行,只能靠任务内部去抢"组锁",但组锁的粒度是任务,一个任务持锁等下游,另一个任务在别的线程里等同一把锁,等待期间线程资源全被占着;如果再不小心搞出任务 A 等任务 B、任务 B 又等任务 A 的情况,就是死锁。
用"组长"模式(同组任务委托给一个专属线程)倒是能解决顺序问题,但一个组一个常驻线程,两百个组就要开两百个线程,线程切换开销和内存占用都压不住。
Bgrp 的做法等价于把组锁从任务层面提升到了调度层面:调度器保证同一个组同一时刻只有一个批次,工作线程永远不需要为"抢组锁"而等待。任务执行期间只需要处理下游的 I/O 等待,线程利用率高得多。不想重复造轮子的同学可以直接用 Celery、Kafka 这种现成方案,但如果你和我一样卡在"既有组内顺序约束、又有下游限流、还要批量化"的中间地带,自己写一个几十行的调度器反而是最划算的。
3. 实现走读:核心代码与三个关键参数
3.1 提交接口与内部数据结构
Bgrp 的核心数据结构非常朴素:一个defaultdict(deque)作为组缓冲区,一个set记录正在执行批次的组,一个Condition做线程间通知,再加一个ThreadPoolExecutor做执行池。下面是最小可运行版本的核心逻辑:
import threading from collections import defaultdict, deque from concurrent.futures import ThreadPoolExecutor class Bgrp: def __init__(self, max_workers=8, batch_size=100, flush_interval=0.2): self.max_workers = max_workers self.batch_size = batch_size self.flush_interval = flush_interval self._pool = ThreadPoolExecutor(max_workers=max_workers) self._bufs = defaultdict(deque) self._running_groups = set() self._cond = threading.Condition() self._closed = False self._cursor = 0 self._dispatch_thread = threading.Thread(target=self._dispatch_loop, daemon=True) self._dispatch_thread.start() def submit(self, group_key, task): with self._cond: self._bufs[group_key].append(task) self._cond.notify() def _dispatch_loop(self): while not self._closed: with self._cond: keys = [k for k, q in self._bufs.items() if k not in self._running_groups and len(q) > 0] if not keys: self._cond.wait(self.flush_interval) continue key = keys[self._cursor % len(keys)] self._cursor += 1 batch = [self._bufs[key].popleft() for _ in range(min(self.batch_size, len(self._bufs[key])))] self._running_groups.add(key) self._pool.submit(self._run_batch, key, batch) def _run_batch(self, key, batch): try: for task in batch: task() finally: with self._cond: self._running_groups.discard(key) self._cond.notify()调用方只需要写一行:bgrp.submit("merchant_1001", lambda: push(merchant_1001, data))。任务可以是任何可调用对象,或者传入参数让内部包装成闭包。
这段代码里最值得注意的点是_run_batch的finally块。无论批次里某个任务抛没抛异常,都必须把组从_running_groups里移除,否则这个组会被永久卡死,后面排队的任务永远没有机会执行。我后来踩过这个坑,一开始只写了正常路径的清理,结果一个任务抛异常,整个组的后继任务全部饿死,业务侧表现为"某个商户的通知突然停发"。
3.2 三个参数怎么定:max_workers、batch_size、flush_interval
Bgrp 对外暴露三个参数,每一个都对应着一个真实的权衡。
max_workers是执行池的线程数。它不决定"有多少个组能同时跑"(那取决于有多少组不在执行中),但决定"最多有多少个批次同时在途"。大多数场景里批次执行的是 I/O 型任务,所以这个值不是按 CPU 核数算的,而是按下游能扛的并发数算的。比如下游总并发上限是 20,批次执行时单批任务连续调用,一个批次占一个并发位,max_workers设成 16 到 18 比较稳妥,留一点余量给下游自己的波动。如果是 CPU 型任务,就老老实实按核心数乘 2 左右来设。
batch_size是每个批次最大任务数。它直接决定攒批的吞吐上限,也决定单批次的延迟上限。批次太大,最后一个任务要多等前面九十九个任务跑完才能轮到,P95 延迟会变差;批次太小,攒批和派发的调度开销占比上升。实测下来,对于 5ms 左右延迟的轻量下游调用,batch_size 在 100 到 500 之间都有不错的收益;如果下游支持真正的批量接口(一次请求传 1000 条),batch_size 可以直接对齐下游的单次上限。
flush_interval是调度器的扫描周期。它决定了"一个只有三五条任务的组,最快多久能被派发"。设成 0.2 秒,意味着小组的任务最坏要等 0.2 秒才开始执行;设成 0.01 秒,调度器每秒醒一百次,空转开销就上来了。实测里 0.1 到 0.2 秒是性价比很高的区间:对拉长任务执行来说,200ms 的起步延迟完全感知不到,但调度线程的 CPU 占用能压到 1% 以下。
3.3 调度循环里的饥饿避免
调度循环里那个_cursor游标,是我处理"组饥饿"问题加的。
如果不加游标,每次都用keys[0],那么第一个组的任务永远被先派发,只要它一直在产生新任务,后面的组就一直被饿着。这在业务上很危险:一个大商户的通知量是普通商户的上百倍,不加干预的话,所有小商户的通知会被无限推迟。
加了游标做轮询之后,每次派发向后移动一位,所有非空组都有机会被选中。这里有一个细节:self._bufs.items()的顺序在 Python 3.7 之后是插入顺序,新出现的组会排在末尾,游标轮询时按len(keys)取模,不会出现新组抢跑的情况。每组一轮最多取一个批次,取完一轮再回来看,整体公平性是够用的。
4. 小实验的实测结果:数据说明设计对不对
4.1 测试环境与压测方法
实验跑在一台 8 核 16G 的 Linux 虚拟机上,Python 3.10。任务数据是模拟生成的:十万条任务,分属 200 个组,每个任务用time.sleep(5ms + random(0-10ms))模拟下游接口延迟,另外用一个计数器模拟下游总并发上限 20,超过就立即抛限流异常,需要重试。
对照组设计了三个方案:方案 A 是我最初的"全局无界队列 + ThreadPoolExecutor(20)";方案 B 是"组级锁 + 固定线程池",任务执行前抢对应组的锁;方案 C 是 Bgrp 本身,max_workers=16,batch_size分别设为 200 和 1000,flush_interval=0.2。每组跑三遍取中位数,记录总耗时、P95 单任务延迟、组内乱序率和运行期间的最大内存。
4.2 对照组的四组数据
| 方案 | 总耗时(秒) | P95 单任务延迟(ms) | 组内乱序率 | 峰值内存 |
|---|---|---|---|---|
| A:全局队列 + 线程池 | 142 | 980 | 5.2% | 210MB |
| B:组级锁 + 线程池 | 96 | 620 | 0% | 180MB |
| C:Bgrp(batch=200) | 61 | 330 | 0% | 120MB |
| C:Bgrp(batch=1000) | 58 | 410 | 0% | 150MB |
方案 A 的乱序率是 5.2%,原因就是两个线程同时取到了同一个组的任务,后取的先完成。这个数字看起来不高,但对通知类业务来说,任何乱序都意味着用户可能收到"欢迎语"和"优惠券到账"顺序颠倒,属于不可接受的事故。
方案 B 把乱序率压到了 0%,但总耗时只比 A 好了三分之一。原因是组级锁让线程在等待锁时空转,锁的竞争和释放本身就是成本。方案 C 的优势是两方面的:调度层保证串行,线程不需要抢锁;批次化减少了调度唤醒的次数,所以总耗时和延迟都比 B 好一截。
4.3 从结果反推的三个结论
第一个结论:批次化是吞吐的第一功臣。A 到 B 只是解决了顺序问题,吞吐提升有限;B 到 C 引入了批,总耗时从 96 秒降到 61 秒,提升约 37%。线程池的价值在执行,批量化的价值在减少执行单位之间的切换成本。
第二个结论:组内串行并不代表吞吐下降。很多人一听"同一个组要串行"就觉得并发白做了,但实测里 200 个组的串行约束并没有拖累整体吞吐,因为不同组的批次在并发执行,单个组的串行只是让"组"这个维度上的并发数为 1,组之间照样是 16 路并行。吞吐只取决于同时在途的批次数量和批次大小,而不是组的数量。
第三个结论:batch_size 不是越大越好。batch=1000 时总耗时略有下降,但 P95 延迟从 330ms 涨到 410ms,峰值内存也从 120MB 涨到 150MB。原因是大批次让排在后面的任务等待更久。实际业务里,延迟和吞吐的平衡点通常在下游接口的单次上限附近,而不是在内存允许的最大值附近。
5. 踩坑记录:这三个细节最容易翻车
5.1 死锁:组内任务等待组外任务
第一次给 Bgrp 接入真实业务时,我遇到了任务全部卡死的情况。排查到最后,问题出在一个商户的通知任务里又调用了bgrp.submit("merchant_1001", retry_task),然后同步等待这个子任务完成。此时父任务正占着这个组的"在途"名额,子任务被调度器判定为"该组正在执行中",永远无法派发。父任务等子任务,子任务等父任务释放名额,直接死锁。
这个坑的根源是把"组串行"和"任务父子关系"混在了一起。Bgrp 的模型假设每个组同一时刻只有一个批次在跑,但批次内部的任务是可以并发或嵌套的,一旦嵌套并且同步等待,就打破了"一个批次必会结束"的前提。
解决方式有两个:一是约定任务内部不允许等待同一个组的其他任务,子任务用"提交后即返回"的方式,等下一轮批次再处理;二是如果业务流程非要同步等待,就不要走 Bgrp,直接在同一批次内部处理完,不做跨批次的等待。我在代码注释里专门加了一行:"禁止在任务内 submit 到同一组并同步等待,否则死锁。"
5.2 分组 Key 倾斜,并行直接失效
第二个坑是数据倾斜。接入的另一个业务里,"组"是客户 ID,但其中一个头部客户的数据量占了总量的八成。Bgrp 的公平轮询只能保证每个组都有机会,不能保证数据多的组不被自己的组内串行约束压住。那个大客户自己的十万条任务必须串行执行,再快的调度也快不起来,整体吞吐被这个单组死死压住。
如果业务允许,可以对热 Key 做拆分:把同一个逻辑组按一定规则拆成多个物理子组,比如客户 ID 加上_0到_7八个后缀,再在下游侧汇总结果。这样做的前提是下游不要求严格有序,或者拆分后的子组之间没有顺序依赖。如果业务必须严格有序,那就没有捷径,只能接受这个瓶颈,或者和业务方商量把强顺序改成最终一致。
这个坑的教训是:任何按 Key 分组的系统,都要提前评估 Key 的分布。做压测的时候不能只测"均匀分布"的数据,一定要把真实的 Key 分布灌进去跑一遍。
5.3 优雅关闭时最容易丢任务
第三个坑出现在服务发版重启的时候。我第一版实现的close()只是_pool.shutdown(wait=True),但线程池关闭只等已在池中的批次跑完,组缓冲区里还没被派发的任务直接就丢了。一次发版,一千多条商户通知悄无声息地消失,第二天业务方来投诉我才发现。
正确做法是三步:先把_closed置为 True,让submit()拒绝新任务;然后让调度线程持续工作直到所有组缓冲区清空、所有在途批次结束;最后再关线程池。代码逻辑大致是这样:
def wait_drain(self, timeout=None): deadline = time.time() + (timeout or 30) while time.time() < deadline: with self._cond: pending = sum(len(q) for q in self._bufs.values()) if pending == 0 and not self._running_groups: return True time.sleep(0.1) return False def close(self, timeout=None): with self._cond: self._closed = True self._cond.notify_all() if not self.wait_drain(timeout): raise TimeoutError("drain timeout") self._pool.shutdown(wait=True)这里有一个容易被忽略的细节:_dispatch_loop里用的是while not self._closed,而wait_drain需要调度线程在关闭后继续派发剩余任务,所以_dispatch_loop的退出条件必须是"已关闭且缓冲区为空",而不是简单的"已关闭"。我在这一版里把循环条件改成了while not self._closed or any(self._bufs.values()):,并在空转时用wait挂起,彻底解决了发版丢任务的问题。
6. 一点后续的想法
Bgrp 这个小实验做完之后,我又陆陆续续给它加了几个小功能:缓冲区积压量的监控指标、批次执行耗时的埋点、以及任务重试次数的限制。这些都不是核心,但对线上运维很有用——至少现在某个组卡住了,我能从指标上立刻看到是哪个组、卡了多久,而不是等业务方来报。
如果让我重新做一遍这个实验,我可能会在开始之前先把"谁等谁"的图先画出来,而不是直接写代码。并发调度的坑,十个里有八个是等待关系没理清导致的,Bgrp 的两层模型把"组内等待"和"组间并行"彻底分开,其实就是在画这张图。对于正在做类似功能的同学,我的个人建议是:先用十分钟把约束条件列全,包括下游限流、顺序要求、失败重试、优雅关闭,再决定用现成方案还是自己写调度。约束列全了,方案基本就出来了。