Kombu 4.1.0:使用Kombu队列时Worker内存占用持续增长(疑似泄漏)
这种内存持续攀升甚至耗尽系统资源的情况我之前碰到过好多次,结合你用的Kombu SimpleQueue + Redis后端的场景,大概率是几个常见的内存泄漏点导致的,咱们一步步来排查解决:
1. 用上下文管理器确保资源自动释放
很多时候内存泄漏是因为队列或连接实例没被正确关闭,尤其是在循环消费消息时反复创建实例但没回收资源。建议用Python的with上下文管理器来创建SimpleQueue,它会自动帮你处理资源的释放:
from kombu import Connection, SimpleQueue # 把连接初始化放在循环外,复用连接 conn = Connection('redis://localhost:6379/') while True: with SimpleQueue(conn) as queue: message = queue.get(block=True) # 这里写你的消息处理逻辑 print(f"处理消息: {message.body}") # 务必确认消息处理完成后ack,避免消息重复被拉取 message.ack()
这样每次消费完消息,队列的相关资源都会被自动释放,不会在内存中堆积。
2. 配置Redis连接池限制连接数
如果没配置连接池,Kombu可能会为每个请求创建新的Redis连接,而旧连接没被及时回收,导致内存里积累大量连接对象。手动配置连接池可以解决这个问题:
from kombu import Connection from redis import ConnectionPool # 配置连接池,限制最大连接数,根据你的业务调整数值 redis_pool = ConnectionPool( host='localhost', port=6379, db=0, max_connections=20 # 按需设置,避免无限制创建连接 ) # 把连接池传给Kombu的Connection conn = Connection( 'redis://localhost:6379/', transport_options={'pool': redis_pool} )
连接池会自动管理连接的复用和回收,避免无限制创建新连接占用内存。
3. 确保消息被正确确认,避免内存堆积
如果处理完消息后没调用message.ack(),Kombu会认为消息没处理完成,可能会把消息留在内存中等待重新处理,时间久了就会占用大量内存。务必在消息处理逻辑完成后调用message.ack(),确认消息已处理完毕。
另外,如果你的消息体很大,处理完后可以显式删除相关变量,触发垃圾回收:
# 处理完大消息后 del large_message_body import gc gc.collect()
不过这是兜底手段,优先确保逻辑里没有不必要的对象引用。
4. 升级Kombu和Redis-py到最新稳定版
某些旧版本的Kombu在Redis传输层存在内存泄漏的BUG,比如连接或队列资源没被正确释放。建议直接升级到最新稳定版:
pip install --upgrade kombu redis
确保redis-py版本和Kombu兼容(推荐用redis-py 4.x以上版本),新版本通常会修复已知的内存泄漏问题。
5. 配置客户端侧的TCP Keepalive
虽然你已经在Redis服务器配置了tcp_keepalive,但Kombu的Redis客户端可能没启用对应配置,导致闲置连接没被及时关闭。可以在创建Connection时添加客户端的keepalive配置:
import socket from kombu import Connection conn = Connection( 'redis://localhost:6379/', transport_options={ 'socket_keepalive': True, 'socket_keepalive_options': { socket.TCP_KEEPIDLE: 60, # 和服务器的tcp_keepalive配置对齐 socket.TCP_KEEPINTVL: 10, socket.TCP_KEEPCNT: 3 } } )
这样客户端和服务器的keepalive机制协同工作,闲置连接会被及时关闭回收。
6. 定期触发Python垃圾回收
如果你的Worker长时间运行,Python的自动垃圾回收可能没及时清理循环引用的对象,可以在循环中定期手动触发垃圾回收,比如每处理100条消息执行一次:
import gc message_count = 0 while True: with SimpleQueue(conn) as queue: msg = queue.get(block=True) # 处理消息 msg.ack() message_count += 1 if message_count % 100 == 0: gc.collect()
这只是临时缓解手段,最好还是找到根本原因,但在排查期间可以先用这个方法控制内存增长。
建议你先从升级Kombu版本和使用上下文管理器管理队列这两点开始排查,这是最常见的解决方向,通常能解决大部分内存泄漏问题。
内容的提问来源于stack exchange,提问作者sradhakrishna

