You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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 = True
    
    将CELERY_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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 01:10:28