Python uWSGI环境下RabbitMQ与日志文件描述符异常排查求助
使用Python + uWSGI部署服务,在非confirm模式下向RabbitMQ发送消息时偶尔出现异常,相关日志及代码如下:
RabbitMQ服务器错误日志
=ERROR REPORT==== 21-Jul-2022::15:23:04 ===
Error on AMQP connection <0.23590.991> (172.198.12.10:59211 -> 10.0.12.1:5672, vhost: 'host', user: 'host', state: running), channel 12848:
operation none caused a connection exception frame_error: "type 91, all octets = <<>>: {frame_too_large,842149168,131064}"
Python应用错误日志
IOError: [Errno 9] Bad file descriptor
Traceback (most recent call last):
File "/usr/lib64/python2.6/logging/init.py", line 800, in emit
self.flush()
File "/usr/lib64/python2.6/logging/init.py", line 762, in flush
self.stream.flush()
日志初始化代码
def init_logger(self): self._logger = logging.getLogger(self._proj) self._logger.setLevel(logging.DEBUG) formatter = logging.Formatter('[%(asctime)s] [%(process)d] [%(levelname)s] %(message)s') if not self._b_stream_init: stream_handler = logging.StreamHandler(sys.stderr) stream_handler.setFormatter(formatter) stream_handler.setLevel(logging.DEBUG) self._logger.addHandler(stream_handler) self._b_stream_init = True ret = self.check_log_name() if ret[0]: return 0 try: log_file_handler = logging.FileHandler(ret[1]) log_file_handler.setFormatter(formatter) log_file_handler.setLevel(CLog.LEVEL_MAP[self._log_level]) self._logger.addHandler(log_file_handler) if self._last_file_handle is not None: self._logger.removeHandler(self._last_file_handle) self._last_file_handle.close() self._last_file_handle = log_file_handler self._last_log_name = ret[1] except: pass def check_log_name(self): if self._log_dir is None or self._log_prefix is None: return True, None log_name_arr = [self._log_dir, self._log_prefix, '_', time.strftime('%Y%m%d_%H'), '.log'] log_name = ''.join(log_name_arr) if self._last_log_name != log_name or not os.path.exists(log_name): return False, log_name else: return True, log_name
RabbitMQ消息发送代码
@trace_report(switch=True) def send(self, key, message, declare=False, expiration=None): if declare: self._declare_exchange() self._declare_queue_with_key(key) if isinstance(expiration, int): expiration = str(expiration) properties = pika.BasicProperties(delivery_mode=2, expiration=expiration) self.channel.basic_publish(exchange=self.exchange, routing_key=key, body=message, properties=properties)
怀疑文件系统将日志内容误写入RabbitMQ的socket fd,当前socket_wait数量达4万但无内核日志,寻求解决方案。
1. 修复多进程下的日志管理问题
uWSGI多进程环境中,自定义日志切分逻辑存在fd复用风险:父进程fork后子进程会继承所有文件描述符,当父/子进程关闭旧日志handle时,释放的fd可能被RabbitMQ的socket占用,导致日志内容写入socket,引发RabbitMQ的frame_too_large错误,同时日志操作时因fd指向socket出现Bad file descriptor。
具体修改:
替换自定义日志切分逻辑,使用官方的TimedRotatingFileHandler实现按小时切分日志,该handler原生支持多进程场景下的文件管理,避免手动处理fd的错误:
def init_logger(self): self._logger = logging.getLogger(self._proj) self._logger.setLevel(logging.DEBUG) formatter = logging.Formatter('[%(asctime)s] [%(process)d] [%(levelname)s] %(message)s') # 保留控制台日志 if not self._b_stream_init: stream_handler = logging.StreamHandler(sys.stderr) stream_handler.setFormatter(formatter) stream_handler.setLevel(logging.DEBUG) self._logger.addHandler(stream_handler) self._b_stream_init = True # 使用TimedRotatingFileHandler按小时切分日志 if self._log_dir and self._log_prefix: log_file_path = os.path.join(self._log_dir, f"{self._log_prefix}.log") # 按小时切分,保留备份文件 file_handler = logging.handlers.TimedRotatingFileHandler( log_file_path, when='H', interval=1, backupCount=24*7, # 保留7天日志 encoding='utf-8' ) file_handler.setFormatter(formatter) file_handler.setLevel(CLog.LEVEL_MAP[self._log_level]) # 避免重复添加handler if not any(isinstance(h, logging.handlers.TimedRotatingFileHandler) for h in self._logger.handlers): self._logger.addHandler(file_handler)
移除原有的check_log_name方法及_last_file_handle、_last_log_name等全局状态变量,避免跨进程的状态干扰。
2. 修复RabbitMQ客户端的多进程兼容性问题
uWSGI多进程环境下,pika的Connection和Channel不能跨进程共享,socket fd继承后多个进程操作同一socket会导致数据混乱:
- 每个uWSGI Worker进程独立初始化RabbitMQ连接和Channel,不要在父进程中提前初始化后让子进程复用;
- 修改发送消息的类,在
__init__或第一次调用send方法时,为当前进程创建独立的连接和channel:
def __init__(self, exchange, ...): self.exchange = exchange # 初始化时不创建连接,延迟到第一次发送时 self._connection = None self._channel = None def _get_channel(self): # 为当前进程创建独立的channel if self._channel is None or not self._channel.is_open: # 创建新连接 credentials = pika.PlainCredentials('host', 'password') parameters = pika.ConnectionParameters('10.0.12.1', 5672, 'host', credentials) self._connection = pika.BlockingConnection(parameters) self._channel = self._connection.channel() # 声明exchange和queue(如果需要) self._declare_exchange() return self._channel @trace_report(switch=True) def send(self, key, message, declare=False, expiration=None): channel = self._get_channel() if declare: self._declare_queue_with_key(key) if isinstance(expiration, int): expiration = str(expiration) properties = pika.BasicProperties(delivery_mode=2, expiration=expiration) channel.basic_publish(exchange=self.exchange, routing_key=key, body=message, properties=properties)
3. 优化uWSGI配置
- 开启
log-master选项,让uWSGI Master进程统一处理日志,避免子进程直接操作日志文件:
[uwsgi] log-master = true
- 限制Worker进程数量,避免过多进程导致fd资源紧张;
- 开启
reload-on-rss等选项,当进程内存占用过高时自动重启,避免fd泄漏积累。
4. 验证与临时恢复
- 使用
lsof -p <uwsgi-worker-pid>查看进程的文件描述符列表,确认日志文件fd与RabbitMQ socket fd是否存在复用情况; - 临时重启uWSGI进程,清理堆积的socket_wait连接;
- 若需临时调整RabbitMQ参数,可修改
frame_max(rabbitmqctl set_vhost_limits -p host '{"frame_max": 1048576}'),但这仅为缓解,需优先解决fd复用根源问题。
内容的提问来源于stack exchange,提问作者lewissi

