【问题标题】:Celery, RabbitMQ messages keep getting sentCelery、RabbitMQ 消息不断发送
【发布时间】:2016-07-13 01:07:52
【问题描述】:

这是设置 - django 项目,其中包含 celery 和一个 CloudAMQP rabbitMQ worker 进行消息代理。

我的 Celery/RabbitMQ 设置:

# RabbitMQ & Celery settings
BROKER_URL = 'ampq://guest:guest@localhost:5672/' # Understandably fake
BROKER_POOL_LIMIT = 1
BROKER_CONNECTION_TIMEOUT = 30
BROKER_HEARTBEAT = 30
CELERY_SEND_EVENTS = False
CELERY_ACCEPT_CONTENT = ['json']
CELERY_TASK_SERIALIZER = 'json'

使用以下命令运行 celery 的 docker 容器:

bash -c 'cd django && celery -A pkm_main worker -E -l info --concurrency=3'

shared_task 定义:

from __future__ import absolute_import

from celery import shared_task

@shared_task
def push_notification(user_id, message):
    logging.critical('Push notifications sent')
    return {'status': 'success'}

当事情发生时我实际上会调用它(我省略了一些代码,因为它似乎不相关):

from notificatons.tasks import push_notification

    def like_this(self, **args):
    # Do like stuff and then do .delay()
    push_notification.delay(media.user.id, request.user.username + ' has liked your item')

所以当它运行时 - 一切看起来都很好 - 输出看起来像这样:

worker_1 | [2016-03-25 09:03:34,888: INFO/MainProcess] Received task: notifications.tasks.push_notification[8443bd88-fa02-4ea4-9bff-8fbec8c91516]
worker_1 | [2016-03-25 09:03:35,333: CRITICAL/Worker-1] Push notifications sent
worker_1 | [2016-03-25 09:03:35,336: INFO/MainProcess] Task notifications.tasks.push_notification[8443bd88-fa02-4ea4-9bff-8fbec8c91516] succeeded in 0.444933412999s: {'status': 'success'}

所以从我收集到的任务已经正确运行和执行,消息应该停止并且 RabbitMQ 应该停止。

但在我的 RabbitMQ 管理中,我看到消息不断发布和传递:

所以我从中收集到的是 RabbitMQ 正在尝试发送某种确认并失败并重试?有没有办法真正关闭这种行为?

热烈欢迎所有帮助和建议。

编辑:忘了提一些重要的事情——直到我调用 push_notification.delay() 之前,消息选项卡是空的,除了每 30 秒来来去去的心跳。只有在我调用 .delay() 之后才会发生这种情况。

编辑 2:CELERYBEAT_SCHEDULE 设置(我试过在有和没有它们的情况下运行 - 没有区别,但添加它们以防万一)

CELERYBEAT_SCHEDULE = {
    "minutely_process_all_notifications": {
        'task': 'transmissions.tasks.process_all_notifications',
        'schedule': crontab(minute='*')
    }
}

编辑 3:添加了查看代码。另外我没有使用 CELERYBEAT_SCHEDULE。我只是将配置保留在代码中以供将来的计划任务使用

from notifications.tasks import push_notification

class MediaLikesView(BaseView):
    def post(self, request, media_id):
        media = self.get_object(media_id)
        data = {}
        data['media'] = media.id
        data['user'] = request.user.id
        serializer = MediaLikeSerializer(data=data)
        if serializer.is_valid():
            like = serializer.save()
            push_notification.delay(media.user.id, request.user.username + ' has liked your item')
            serializer = MediaGetLikeSerializer(like)
            return self.get_mocked_pagination_response(status=status.HTTP_204_NO_CONTENT)
        return self.get_mocked_pagination_response(serializer.errors, status=status.HTTP_400_BAD_REQUEST)

【问题讨论】:

  • 能否请您发布您调用方法like_this 的代码?以及CELERYBEAT_SCHEDULE 的值(如果有)?
  • 嘿。我已经添加了时间表。该方法本身位于 Django Rest Framework 视图类的 post 方法中。不知道这是否有帮助,但我也可以把它扔在那里。
  • 方法process_all_notifications在哪里?我的意思是,它应该在transmissions.tasks,但是我没有看到这个代码,你能不能也加一下?
  • 请同时添加视图代码
  • 没有这种方法。如果我想添加 CELERYBEAT 并执行 crons,这只是未来的占位符设置。

标签: python django rabbitmq celery


【解决方案1】:

这是芹菜的混合和八卦。通过将--without-gossip --without-mingle --without-heartbeat 添加到命令行参数来禁用。

当你在命令行上禁用心跳时不要忘记设置BROKER_HEARTBEAT = None,否则你会在30s后断开连接。通常依赖 TCP keepalive 比 AMQP 心跳更好,或者更糟糕的是,Celery 自己的心跳。

【讨论】:

  • 卡尔,愿你碰巧相信的任何饮食都能保佑你和你的全家。那成功了。非常感谢你。我要去阅读八卦和混合的作用。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-12-31
  • 1970-01-01
相关资源
最近更新 更多