RabbitMQ达队列阈值时意外关闭Celery TCP连接问题求助
问题现象
采用RabbitMQ 3.8.2 + Erlang 22.2.7 + django-celery-rabbitmq架构,消费任务时出现以下异常:
- 当队列消息数达到1200条(单条消息含1000-1500字符)时,RabbitMQ开始关闭AMQP连接,日志输出"client unexpectedly closed TCP connection"等错误;
- 随后django-celery消费者从队列列表中消失,Celery Pod完成任务后无法ack消息,抛出
ConnectionResetError; - 若单条消息缩小至50字符,触发阈值升至4000-5000条;
- 服务器CPU、内存资源充足,已尝试调整Celery预取数、心跳等参数,问题未解决。
RabbitMQ日志
... 2022-11-01 09:35:25.327 [info] <0.20608.9> 接受AMQP连接 <0.20608.9> (185.121.83.107:60447 -> 185.121.83.116:5672) 2022-11-01 09:35:25.483 [info] <0.20608.9> 连接 <0.20608.9> (185.121.83.107:60447 -> 185.121.83.116:5672): 用户 'rabbit_admin' 认证通过并获得虚拟主机 '/' 的访问权限 ... 2022-11-01 09:36:59.129 [warning] <0.19994.9> 关闭AMQP连接 <0.19994.9> (185.121.83.108:36149 -> 185.121.83.116:5672, 虚拟主机: '/', 用户: 'rabbit_admin'): 客户端意外关闭TCP连接 ... [error] <0.11162.9> 关闭AMQP连接 <0.11162.9> (185.121.83.108:57631 -> 185.121.83.116:5672): {writer,send_failed,{error,enotconn}} ... 2022-11-01 09:35:48.256 [error] <0.20201.9> 关闭AMQP连接 <0.20201.9> (185.121.83.108:50058 -> 185.121.83.116:5672): {inet_error,enotconn} ...
Celery报错日志
ERROR: [2022-11-01 09:20:23] /usr/src/app/project/celery.py:114 handle_message 处理Rabbit任务时出错: [Errno 104] 连接被对等方重置 Traceback (最近的调用最后): 文件 "/usr/local/lib/python3.10/site-packages/amqp/connection.py", 第514行, 在channel中 return self.channels[channel_id] KeyError: None 处理上述异常时,发生了另一个异常: Traceback (最近的调用最后): 文件 "/usr/src/app/project/celery.py", 第76行, 在handle_message中 message.ack() 文件 "/usr/local/lib/python3.10/site-packages/kombu/message.py", 第125行, 在ack中 self.channel.basic_ack(self.delivery_tag, multiple=multiple) 文件 "/usr/local/lib/python3.10/site-packages/amqp/channel.py", 第1407行, 在basic_ack中 return self.send_method( 文件 "/usr/local/lib/python3.10/site-packages/amqp/abstract_channel.py", 第70行, 在send_method中 conn.frame_writer(1, self.channel_id, sig, args, content) 文件 "/usr/local/lib/python3.10/site-packages/amqp/method_framing.py", 第186行, 在write_frame中 write(buffer_store.view[:offset]) 文件 "/usr/local/lib/python3.10/site-packages/amqp/transport.py", 第347行, 在write中 self._write(s) ConnectionResetError: [Errno 104] 连接被对等方重置
Celery客户端配置
CELERY_TASK_IGNORE_RESULT = True CELERY_RESULT_BACKEND = 'django-db' CELERY_CACHE_BACKEND = 'django-cache' CELERY_SEND_EVENTS = False CELERY_BROKER_POOL_LIMIT = 30 CELERY_BROKER_HEARTBEAT = 30 CELERY_BROKER_CONNECTION_TIMEOUT = 600 CELERY_PREFETCH_MULTIPLIER = 1 CELERY_SEND_EVENTS = False CELERY_WORKER_CONCURRENCY = 1 CELERY_TASK_ACKS_LATE = True
Celery启动命令
当前启动命令:
celery -A project.celery worker -l info -f /var/log/celery/celery.log -Ofair
已尝试的调整命令:
celery -A project.celery worker -l info -f /var/log/celery/celery.log --without-heartbeat --without-gossip --without-mingle
celery -A project.celery worker -l info -f /var/log/celery/celery.log --prefetch-multiplier=1 --pool=solo
排查及解决建议
- 检查RabbitMQ连接限制参数
用rabbitmqctl list_parameters查看connection_max和channel_max,若设置了连接/通道上限,调整为更大值或取消限制,避免消息量达标后触发连接关闭。 - 调整RabbitMQ资源阈值
用rabbitmqctl environment查看vm_memory_high_watermark,建议调整为0.7(70%内存使用率),避免内存保护机制误触发;同时确认磁盘空闲空间满足disk_free_limit要求,防止因磁盘阈值触发连接关闭。 - 优化Celery连接配置
降低CELERY_BROKER_POOL_LIMIT到5-10(当前30在单并发场景下易过载);添加自动重连配置:
将CELERY_BROKER_CONNECTION_RETRY = True CELERY_BROKER_CONNECTION_RETRY_ON_STARTUP = TrueCELERY_BROKER_HEARTBEAT调整为60,与RabbitMQ默认心跳值匹配,避免心跳超时误判连接失效。 - 排查网络层面限制
检查防火墙、负载均衡是否有连接超时或流量限制,临时关闭防火墙测试;调整系统TCP参数(如tcp_keepalive_time、tcp_keepalive_intvl),避免操作系统过早回收连接。 - 升级组件版本
RabbitMQ 3.8.2和Erlang 22.2.7存在已知连接稳定性问题,建议升级到RabbitMQ 3.8.x最新补丁版(如3.8.35)+ Erlang 23.x/24.x(注意版本兼容);同时升级Celery、kombu到最新稳定版,修复旧版本AMQP连接处理Bug。 - 调整消息处理策略
若业务允许,拆分大消息为小消息,降低单条消息负载;启用任务批量处理,减少连接交互开销。
内容的提问来源于stack exchange,提问作者Kaiser
相关产品推荐
相关产品推荐

