Celery add_periodic_task阻止Django在uwsgi环境中运行

Tim*_*Tim 5 python django celery celery-task celerybeat

我编写了一个模块,根据项目设置中的字典列表(通过导入django.conf.settings)动态添加定期芹菜任务.我这样做是使用一个函数add_tasks来调度一个函数,该函数使用uuid设置中给出的特定函数进行调用:

def add_tasks(celery):
    for new_task in settings.NEW_TASKS:
        celery.add_periodic_task(
            new_task['interval'],
            my_task.s(new_task['uuid']),
            name='My Task %s' % new_task['uuid'],
        )
Run Code Online (Sandbox Code Playgroud)

这里建议我使用on_after_configure.connect信号来调用我的函数celery.py:

app = Celery('my_app')

@app.on_after_configure.connect
def setup_periodic_tasks(celery, **kwargs):
    from add_tasks_module import add_tasks
    add_tasks(celery)
Run Code Online (Sandbox Code Playgroud)

这个设置适用于两者celery beat,celery worker但在我用于uwsgi服务我的django应用程序的地方中断了我的设置.Uwsgi运行顺利,直到视图代码第一次使用celery的.delay()方法发送任务.在这一点上,似乎芹菜被初始化,uwsgi但在上面的代码中永远阻止.如果我从命令行手动运行它然后在它阻塞时中断,我得到以下(缩短的)堆栈跟踪:

Traceback (most recent call last):
  File "/usr/local/lib/python3.6/site-packages/kombu/utils/objects.py", line 42, in __get__
    return obj.__dict__[self.__name__]
KeyError: 'tasks'

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/usr/local/lib/python3.6/site-packages/kombu/utils/objects.py", line 42, in __get__
    return obj.__dict__[self.__name__]
KeyError: 'data'

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/usr/local/lib/python3.6/site-packages/kombu/utils/objects.py", line 42, in __get__
    return obj.__dict__[self.__name__]
KeyError: 'tasks'

During handling of the above exception, another exception occurred:
Traceback (most recent call last):

  (SHORTENED HERE. Just contained the trace from the console through my call to this function)

  File "/opt/my_app/add_tasks_module/__init__.py", line 42, in add_tasks
    my_task.s(new_task['uuid']),
  File "/usr/local/lib/python3.6/site-packages/celery/local.py", line 146, in __getattr__
    return getattr(self._get_current_object(), name)
  File "/usr/local/lib/python3.6/site-packages/celery/local.py", line 109, in _get_current_object
    return loc(*self.__args, **self.__kwargs)
  File "/usr/local/lib/python3.6/site-packages/celery/app/__init__.py", line 72, in task_by_cons
    return app.tasks[
  File "/usr/local/lib/python3.6/site-packages/kombu/utils/objects.py", line 44, in __get__
    value = obj.__dict__[self.__name__] = self.__get(obj)
  File "/usr/local/lib/python3.6/site-packages/celery/app/base.py", line 1228, in tasks
    self.finalize(auto=True)
  File "/usr/local/lib/python3.6/site-packages/celery/app/base.py", line 507, in finalize
    with self._finalize_mutex:
Run Code Online (Sandbox Code Playgroud)

获取互斥锁似乎存在问题.

目前我正在使用一种解决方法来检测是否sys.argv[0]包含uwsgi然后不添加周期性任务,因为只beat需要任务,但我想了解这里出了什么问题来更永久地解决问题.

这个问题可能与使用uwsgi多线程或多处理有关,其中一个线程/进程持有另一个需要的互斥锁吗?

我很感激可以帮助我解决问题的任何提示.谢谢.

我正在使用:Django 1.11.7和Celery 4.1.0

编辑1

我已经为这个问题创建了一个最小的设置:

celery.py:

import os
from celery import Celery
from django.conf import settings
from myapp.tasks import my_task

# set the default Django settings module for the 'celery' program.
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'my_app.settings')

app = Celery('my_app')

@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    sender.add_periodic_task(
        60,
        my_task.s(),
        name='Testtask'
    )

app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)
Run Code Online (Sandbox Code Playgroud)

tasks.py:

from celery import shared_task
@shared_task()
def my_task():
    print('ran')
Run Code Online (Sandbox Code Playgroud)

确保CELERY_TASK_ALWAYS_EAGER = False,并且您有一个有效的消息队列.

跑:

./manage.py shell -c 'from myapp.tasks import my_task; my_task.delay()'
Run Code Online (Sandbox Code Playgroud)

在中断前等待大约10秒钟以查看上述错误.

Tim*_*Tim 1

所以,我发现@shared_task装饰器造成了问题。当我在信号调用的函数中声明任务时,我可以避免这个问题,如下所示:

def add_tasks(celery):
    @celery.task
    def my_task(uuid):
        print(uuid)

    for new_task in settings.NEW_TASKS:
        celery.add_periodic_task(
            new_task['interval'],
            my_task.s(new_task['uuid']),
            name='My Task %s' % new_task['uuid'],
        )
Run Code Online (Sandbox Code Playgroud)

这个解决方案实际上对我有用,但我还有一个问题:我在可插入应用程序中使用此代码,因此我无法直接访问信号处理程序之外的 celery 应用程序,但也希望能够调用my_task其他代码中的函数。通过在函数内定义它,它在函数外部不可用,因此我无法将其导入到其他任何地方。

我可能可以通过在信号函数之外定义任务函数来解决这个问题,并在此处和tasks.py. 我想知道除了装饰器之外是否还有一个装饰器@shared_task可以在tasks.py不会产生问题的情况下使用。

目前最好的解决方案可能是:

task_app.__init__.py:

def my_task(uuid):
    # do stuff
    print(uuid)

def add_tasks(celery):
    celery_my_task = celery.task(my_task)
    for new_task in settings.NEW_TASKS:
        celery.add_periodic_task(
            new_task['interval'],
            celery_my_task(new_task['uuid']),
            name='My Task %s' % new_task['uuid'],
        )
Run Code Online (Sandbox Code Playgroud)

task_app.tasks.py:

from celery import shared_task
from task_app import my_task
shared_my_task = shared_task(my_task)
Run Code Online (Sandbox Code Playgroud)

myapp.celery.py:

import os
from celery import Celery
from django.conf import settings


# set the default Django settings module for the 'celery' program.
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'my_app.settings')

app = Celery('my_app')

@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    from task_app import add_tasks
    add_tasks(sender)


app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)
Run Code Online (Sandbox Code Playgroud)