前段时间处理一个数据批量更新任务,需要从 MySQL 表中读取大约 300 万条记录,经过一系列逻辑处理(调用外部 API 更新标签)后写回。为了控制资源,我使用了 ThreadPoolExecutor 配合生成器逐条产出数据,期望做到“边读边处理”,避免一次性加载全部结果集。
然而上线后观察监控,发现 Python 进程的内存占用随着迭代不断攀升,最终几乎吃满机器内存,导致 OOM 风险。最初直觉是 yield 没有真正实现流式,或者数据库驱动缓存了全部结果,但深入排查后发现事情没那么简单。
初始代码结构
简化后的核心代码如下:
from mysql.connector.pooling import MySQLConnectionPool
dbpool = MySQLConnectionPool()
def fetch(sql, params=None):
with dbpool.get_connection() as conn, conn.cursor(dictionary=True) as cursor:
cursor.execute(sql, params)
yield from cursor
with ThreadPoolExecutor(max_workers=10) as executor:
for idx, item in enumerate(fetch(SQL_SELECT)):
executor.submit(update_label, idx, item)
fetch 是一个生成器函数,使用 yield from 逐行产出。主循环每拿到一行就提交一个任务给线程池,update_label 负责具体处理(包含少量 I/O 和计算)。
与此同时, 经官方文档查询, mysqlconnection-cursor , 参数buffered默认为False, 意味着数据库游标默认使用SSCursor,此时在 cursor.execute() 执行后不会将全部结果集从 MySQL Server 发送回来,而是在客户端迭代时才逐条传输。这能极大降低客户端内存。
经过上面两个处理,逻辑看起来清晰,内存理应稳定,但是事实并不如此。
内存持续增长的真正元凶
通过 memory_profiler 和 tracemalloc 定位,发现内存增长呈现“阶梯式”:主循环迭代速度极快(每秒产出数千行),而工作线程处理速度相对较慢(涉及外部请求,平均耗时 ~50ms)。导致 ThreadPoolExecutor 内部的无界任务队列(_work_queue)迅速堆积,未执行的任务及其携带的参数(item 字典)始终被引用,无法被垃圾回收。
终于查到元凶,原来是线程池队列的问题。
尝试修改队列为有界阻塞——一个危险的念头
排查过程中曾考虑将 executor._work_queue 替换为 Queue(maxsize=N),利用生产者阻塞来反向压制主循环,避免队列无限膨胀。但深入阅读 ThreadPoolExecutor 源码后发现,shutdown() 方法依赖向队列中放入 None 作为哨兵,通知工作线程退出。若队列为有界且已满,shutdown 将阻塞等待空位,而工作线程可能因为处理耗时任务无法及时消费,从而导致死锁。这种修改私有属性的方式风险极高,且不被官方支持,果断放弃。
最终采用的解决方案
方案一:信号量限流(生产环境首选)
不触碰线程池内部队列,通过 Semaphore 控制同时处于“已提交但未完成”状态的任务总数,形成自然背压(Backpressure)。当排队任务达到阈值时,主循环在 sem.acquire() 处阻塞,直到有任务完成释放许可证。
from threading import Semaphore
sem = Semaphore(100) # 限制未完成任务数不超过 100
with ThreadPoolExecutor(max_workers = 15) as executor:
for idx, item in enumerate(fetch(SQL_SELECT)):
sem.acquire()
future = executor.submit(update_label, idx, item)
future.add_done_callback(lambda f: sem.release())
实际测试中,内存从最初的无止境增加降至不足 200MB,效果非常拔群。
方案二:自定义消费者队列(完全控制流)
对于更复杂的场景,也可以弃用 ThreadPoolExecutor 的内建队列,自行维护一个阻塞有界队列,并启动固定数量的工作线程。主线程作为生产者,队列满时阻塞,线程退出信号由外部控制。这种模式虽然代码量稍多,但逻辑清晰,适合需要精细调优的场景。
也是果断放弃。
最终生产配置
最终稳定运行的版本采用了“yield + 服务端游标 + 信号量限流”的组合:
- 数据库连接池最大连接数 20。
- 信号量限制未完成任务数为
max_workers * 2(即 20)。 - 每个
update_label设置超时和重试,避免单个任务卡死导致信号量泄漏。 - 线程池
max_workers=10,经压测后 CPU 和内存均平稳。
300 万条数据全量处理耗时约 3 小时,内存峰值控制在 250MB 以内,较初版下降近 90%。
总结
- 生成器不等于流式加载:要使用 SSCursor 才是真正意义上的逐行流式拉取。
- 线程池的无界队列是隐藏的“内存黑洞”:当生产者速率远超消费者时,任务队列会无限制膨胀,必须引入限流机制。
- 慎改私有变量:
ThreadPoolExecutor的_work_queue虽然可赋值,但其内部 shutdown 逻辑强依赖队列行为,强行替换为有界队列极易引发死锁。 - 监控先行:借助
tracemalloc或memory-profiler定位内存分配来源,往往比凭直觉猜效率更高。