【问题标题】:How to process a Celery task on a specific vhost?如何在特定的虚拟主机上处理 Celery 任务?
【发布时间】:2018-04-13 10:32:12
【问题描述】:

我有一个 Celery 任务,例如:

from celery.task import task
from django.conf import settings
from base.tasks import BaseTask

@task(name="throw_exception", base=BaseTask)
def print_value(*args, **kwargs):
    print('BROKER_URL:', settings.BROKER_URL)

我正在我的 virtualenv 中运行一个 Celery worker,例如:

celery worker -A myproject -l info

工人展示:

Connected to amqp://guest:**@127.0.0.1:5672/myapp

当我从 Django shell 启动我的任务时:

>>> from django.conf import settings
>>> settings.BROKER_URL
'amqp://guest:**@127.0.0.1:5672/myapp'
>>> from myapp.tasks import print_value
>>> print_value.delay()

我从未在工作人员的日志中看到执行的任务。

但是,如果我改为使用带有默认“/”虚拟主机的 BROKER_URL,那么它会立即执行所有待处理的任务,这意味着我对 print_value.delay() 的所有调用都将其发送到错误的虚拟主机,即使设置了正确的 BROKER_URL。我做错了什么?

编辑:问题似乎是 Celery 没有一致的 @task 装饰器,并且通过使用错误的装饰器,您将任务与代理设置断开连接。所以本质上,我的所有任务都配置为使用默认代理,而不是我的设置中定义的代理。旧文档说要使用from celery.task import task,但新文档...并没有真正指定,似乎暗示您应该使用celery.py 文件中定义的app 实例,例如@app.task。问题是我所有的任务都在单独的tasks.py 文件中,他们无法访问app 实例。如果我将一个任务复制到我的celery.py 并使用@app.task 装饰器,那么它会使用正确的虚拟主机并按预期工作,但很明显,这不是一个实际的解决方案,因为我必须复制几十个函数到这个文件中。如何正确解决此问题?

【问题讨论】:

  • 通过sudo rabbitmqctl add_vhost {vhost_name}添加虚拟主机并通过sudo rabbitmqctl set_permissions -p {vhost_name} {username} ".*" ".*" ".*"将用户权限设置为虚拟主机?
  • @Ykh,是的,这不是 rabbitmq 的问题。
  • 你可以尝试在views.py中通过请求方法调用任务,而不是在python shell中,不确定python shell是否使用与你的virtualenv相同的环境。

标签: python django celery


【解决方案1】:

我现在使用 Django + (Celery + RabbitMQ) 有同样的问题。我的解决方案是,

CELERY_BROKER_URL=amqp://<user>:<password>@localhost:5672/<vhost>

这是来自 RabbitMQ.com 的详细确认 > 客户端文档 > RabbitMQ URI Specification


物有所值...

...Celery + RabbitMQ 有很多功能。我正在使用rabbitmqctl list_vhosts-- 查看 RabbitMQ,但我没有看到我的虚拟主机。 WTF? 最后,我意识到我在本地开发服务器上配置supervisord 太早了。从 CLI 启动 Celery 会给出一堆反馈,supervisord 会将这些反馈放在我眼皮底下的某个地方,例如:

[2021-02-19 18:26:52,803: WARNING/MainProcess] (0, 0): (403) ACCESS_REFUSED - Login was refused using authentication mechanism AMQPLAIN. For details see the broker logfile.

AMQP 立即让我想起了连接字符串。繁荣。这是你的vhost

【讨论】:

  • 仅供参考,文档为您提供了一个带有 pyamqp://USER:PASSWORD@HOSTNAME:PORT/ 的示例 URL,我认为尾随 / 是故意表示 rabbitmq vhost 是 /(rabbitmq 始终具有默认 vhost @987654332 @)。
【解决方案2】:

在挖掘了 Celery 的代码之后,我能找到设置当前应用程序的唯一方法是调用 celery._state._set_current_app(app)。显然,这是一种内部方法,并非旨在以这种方式使用,但我找不到任何其他方法将我的自定义应用程序实例显式设置为“当前”应用程序。我原以为这应该自动完成,特别是因为我的代码直接取自教程,所以要么文档不完整,要么这是一个错误。

无论如何,工作 celery 文件看起来像:

from __future__ import absolute_import, print_function
import os
import sys

from celery import Celery
from celery._state import _set_current_app
import django

app = Celery('myproject')

app.config_from_object('django.conf:settings', namespace='CELERY')
_set_current_app(app)

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'myproject.settings.settings')
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), '../myproject')))
django.setup()
from django.conf import settings
app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)

这导致我所有 tasks.py 文件中的所有 @task 装饰器正确访问我的自定义 Celery 实例。

【讨论】:

    【解决方案3】:

    只需提供一个带有不同虚拟主机的演示使用 djcelery 即可。

    在您的项目文件夹中,__init__.py:

    from __future__ import absolute_import
    # This will make sure the app is always imported when
    # Django starts so that shared_task will use this app.
    from .celery import app as celery_app
    

    celery.py,用你自己的项目标签替换SchoolMS

    from __future__ import absolute_import
    import os
    from celery import Celery, platforms
    from django.conf import settings
    
    # set the default Django settings module for the 'celery' program.  
    os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'SchoolMS.settings')
    
    app = Celery('SchoolMS')
    platforms.C_FORCE_ROOT = True
    
    # Using a string here means the worker will not have to  
    # pickle the object when using Windows.  
    app.config_from_object('django.conf:settings')
    app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)
    

    settings.py:

    BROKER_URL = 'amqp://schoolms:schoolms@localhost:5672/schoolms'
    CELERY_TIMEZONE = 'Asia/Shanghai'
    CELERYBEAT_SCHEDULER = 'djcelery.schedulers.DatabaseScheduler'
    CELERYBEAT_SCHEDULE = {
    }
    

    user/tasks.py:

    from celery import task
    
    from django.conf import settings
    
    
    @task
    def send_tel_verify(tel_verify_id):
        try:
            tel_verify = TelVerify.objects.get(id=tel_verify_id)
            try:
                send_sms(tel_verify.tel, 'xxxx')
                return ''success'
            except SmsError as e:
                return 'error'
        except ObjectDoesNotExist:
            return 'not found'
    

    user/views.py

    send_tel_verify.delay(tel_verify.id)
    

    【讨论】:

      猜你喜欢
      • 2020-12-15
      • 2020-01-17
      • 2012-11-27
      • 2018-06-03
      • 2020-07-17
      • 1970-01-01
      • 1970-01-01
      • 2012-09-22
      • 2018-06-20
      相关资源
      最近更新 更多