浏览代码

connection_diagnostics now supports carrot versions >= 0.4.0

Ask Solem 16 年之前
父节点
当前提交
434a1350b1
共有 1 个文件被更改,包括 6 次插入1 次删除
  1. 6 1
      celery/worker.py

+ 6 - 1
celery/worker.py

@@ -250,7 +250,12 @@ class TaskDaemon(object):
     def connection_diagnostics(self):
         """Diagnose the AMQP connection, and reset connection if
         necessary."""
-        if not self.task_consumer.channel.connection:
+        if hasattr(self.task_consumer.backend):
+            connection = self.task_consumer.backend.channel.connection
+        else:
+            connection = self.task_consumer.channel.connection
+
+        if not connection:
             self.logger.info(
                     "AMQP Connection has died, restoring connection.")
             self.reset_connection()