高负载环境下Python进程池Redis连接错误排查求助
问题场景与故障
环境配置
- 操作系统:CentOS7.8
- 服务器CPU:100核
- Redis版本:7.0.11
- Python版本:3.8.10
代码逻辑
通过300个进程的进程池处理请求,每个进程仅初始化一次Redis客户端并复用,核心逻辑:
- 用Redis的
getset命令检查唯一ID对应的key - 若key不存在,存入值并设置60秒过期
- 若key已存在,执行高CPU负载计算后删除key
伪代码如下:
import redis from concurrent.futures import ProcessPoolExecutor redis: redis.Redis def init_process(): global redis redis = redis.Redis(host="127.0.0.1", port=6379, decode_response=True) def recieved_request(unique_id): x = redis.getset(unique_id, "hello world") if x is not None: heavy_cpu_load() redis.delete(unique_id) else: redis.expire(unique_id, 60) pool = ProcessPoolExecutor(max_workers=300, initializer=init_process) def on_new_request(): pool.submit(recieved_request)
故障现象
低流量时运行正常,高流量下频繁出现三类错误:
error 99 - connecting to cannot assign requested addresserror while reading from 127.0.0.1:6379: (104, connection is reset by peer)connection is closed by server
已尝试无效操作:调整连接超时快速关闭空闲连接、将maxconnection设为50000
补充信息:移除高CPU计算逻辑后无错误,无法迁移Redis到其他机器
根因分析
核心问题是高CPU计算导致Redis连接长时间闲置:
- 进程执行
heavy_cpu_load()时,Redis客户端连接处于无交互状态,若Redis配置了连接超时(或TCP层面keepalive未生效),服务器会主动切断连接; - 连接被切断后,客户端进程尝试重建连接,高流量下大量连接进入TIME_WAIT状态,耗尽本地可用端口,触发
cannot assign requested address错误; - 300个进程远超100核CPU的处理能力,进程上下文切换频繁,进一步拉长CPU计算耗时,加剧连接闲置断开的概率。
修复方案
1. 避免连接在计算期间闲置
在执行高CPU计算前主动关闭Redis连接,计算完成后重新初始化:
def recieved_request(unique_id): global redis x = redis.getset(unique_id, "hello world") if x is not None: # 计算前关闭闲置连接 redis.close() heavy_cpu_load() # 计算完成后重建连接 redis = redis.Redis(host="127.0.0.1", port=6379, decode_response=True) redis.delete(unique_id) else: redis.expire(unique_id, 60)
2. 调整TCP参数缓解端口耗尽
编辑/etc/sysctl.conf添加以下配置:
# 允许复用TIME_WAIT状态的端口 net.ipv4.tcp_tw_reuse = 1 # 缩短TIME_WAIT超时时间 net.ipv4.tcp_fin_timeout = 30 # 扩大本地可用端口范围 net.ipv4.ip_local_port_range = 1024 65535
执行sysctl -p使配置生效。
3. 优化进程池大小
进程数远超CPU核心数会加剧上下文切换,建议将max_workers调整为CPU核心数的11.5倍(100150):
pool = ProcessPoolExecutor(max_workers=150, initializer=init_process)
4. 用连接池实现自动重连
改用Redis连接池管理连接,自动处理连接断开后的重连逻辑:
import redis from concurrent.futures import ProcessPoolExecutor from redis import ConnectionPool pool: ConnectionPool def init_process(): global pool pool = ConnectionPool(host="127.0.0.1", port=6379, decode_response=True) def recieved_request(unique_id): redis_cli = redis.Redis(connection_pool=pool) try: x = redis_cli.getset(unique_id, "hello world") if x is not None: heavy_cpu_load() redis_cli.delete(unique_id) else: redis_cli.expire(unique_id, 60) except redis.ConnectionError: # 捕获连接错误后重试 redis_cli = redis.Redis(connection_pool=pool) x = redis_cli.getset(unique_id, "hello world") if x is not None: heavy_cpu_load() redis_cli.delete(unique_id) else: redis_cli.expire(unique_id, 60) pool = ProcessPoolExecutor(max_workers=150, initializer=init_process) def on_new_request(): pool.submit(recieved_request)
内容的提问来源于stack exchange,提问作者Tornado
相关产品推荐
相关产品推荐

