ARTICLE DETAIL

资讯详情

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

Bgrp:分组并发调度框架的设计与实战,解决组内串行与组间并行

Bgrp:分组并发调度框架的设计与实战,解决组内串行与组间并行 几个月前我接了一个内部系统的改造任务需求看起来非常简单上游每天会推过来几十万条带分组标识的任务数据下游接口对并发有硬性限制——同一个分组内的请求必须串行处理不同分组之间可以并行。我第一版直接用线程池硬写结果跑了一周就发现问题不断某个分组处理得慢整个池子跟着被拖垮下游限流一触发所有请求一起遭殃。于是我用几天时间写了一个小工具取名 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_workers8, batch_size100, flush_interval0.2): self.max_workers max_workers self.batch_size batch_size self.flush_interval flush_interval self._pool ThreadPoolExecutor(max_workersmax_workers) self._bufs defaultdict(deque) self._running_groups set() self._cond threading.Condition() self._closed False self._cursor 0 self._dispatch_thread threading.Thread(targetself._dispatch_loop, daemonTrue) 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_intervalBgrp 对外暴露三个参数每一个都对应着一个真实的权衡。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_workers16batch_size分别设为 200 和 1000flush_interval0.2。每组跑三遍取中位数记录总耗时、P95 单任务延迟、组内乱序率和运行期间的最大内存。4.2 对照组的四组数据方案总耗时秒P95 单任务延迟ms组内乱序率峰值内存A全局队列 线程池1429805.2%210MBB组级锁 线程池966200%180MBCBgrpbatch200613300%120MBCBgrpbatch1000584100%150MB方案 A 的乱序率是 5.2%原因就是两个线程同时取到了同一个组的任务后取的先完成。这个数字看起来不高但对通知类业务来说任何乱序都意味着用户可能收到欢迎语和优惠券到账顺序颠倒属于不可接受的事故。方案 B 把乱序率压到了 0%但总耗时只比 A 好了三分之一。原因是组级锁让线程在等待锁时空转锁的竞争和释放本身就是成本。方案 C 的优势是两方面的调度层保证串行线程不需要抢锁批次化减少了调度唤醒的次数所以总耗时和延迟都比 B 好一截。4.3 从结果反推的三个结论第一个结论批次化是吞吐的第一功臣。A 到 B 只是解决了顺序问题吞吐提升有限B 到 C 引入了批总耗时从 96 秒降到 61 秒提升约 37%。线程池的价值在执行批量化的价值在减少执行单位之间的切换成本。第二个结论组内串行并不代表吞吐下降。很多人一听同一个组要串行就觉得并发白做了但实测里 200 个组的串行约束并没有拖累整体吞吐因为不同组的批次在并发执行单个组的串行只是让组这个维度上的并发数为 1组之间照样是 16 路并行。吞吐只取决于同时在途的批次数量和批次大小而不是组的数量。第三个结论batch_size 不是越大越好。batch1000 时总耗时略有下降但 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(waitTrue)但线程池关闭只等已在池中的批次跑完组缓冲区里还没被派发的任务直接就丢了。一次发版一千多条商户通知悄无声息地消失第二天业务方来投诉我才发现。正确做法是三步先把_closed置为 True让submit()拒绝新任务然后让调度线程持续工作直到所有组缓冲区清空、所有在途批次结束最后再关线程池。代码逻辑大致是这样def wait_drain(self, timeoutNone): 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, timeoutNone): with self._cond: self._closed True self._cond.notify_all() if not self.wait_drain(timeout): raise TimeoutError(drain timeout) self._pool.shutdown(waitTrue)这里有一个容易被忽略的细节_dispatch_loop里用的是while not self._closed而wait_drain需要调度线程在关闭后继续派发剩余任务所以_dispatch_loop的退出条件必须是已关闭且缓冲区为空而不是简单的已关闭。我在这一版里把循环条件改成了while not self._closed or any(self._bufs.values()):并在空转时用wait挂起彻底解决了发版丢任务的问题。6. 一点后续的想法Bgrp 这个小实验做完之后我又陆陆续续给它加了几个小功能缓冲区积压量的监控指标、批次执行耗时的埋点、以及任务重试次数的限制。这些都不是核心但对线上运维很有用——至少现在某个组卡住了我能从指标上立刻看到是哪个组、卡了多久而不是等业务方来报。如果让我重新做一遍这个实验我可能会在开始之前先把谁等谁的图先画出来而不是直接写代码。并发调度的坑十个里有八个是等待关系没理清导致的Bgrp 的两层模型把组内等待和组间并行彻底分开其实就是在画这张图。对于正在做类似功能的同学我的个人建议是先用十分钟把约束条件列全包括下游限流、顺序要求、失败重试、优雅关闭再决定用现成方案还是自己写调度。约束列全了方案基本就出来了。
返回列表