【问题标题】:Connect new celery periodic task in django在 django 中连接新的 celery 周期性任务
【发布时间】:2017-04-28 09:26:18
【问题描述】:

这不是一个问题,而是帮助那些会发现 celery 4.0.1 文档中描述的周期性任务声明很难集成到 django 中的人: http://docs.celeryproject.org/en/latest/userguide/periodic-tasks.html#entries

复制粘贴 celery 配置文件main_app/celery.py:

from celery import Celery
from celery.schedules import crontab

app = Celery()

@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    # Calls test('hello') every 10 seconds.
    sender.add_periodic_task(10.0, test.s('hello'), name='add every 10')

    # Calls test('world') every 30 seconds
    sender.add_periodic_task(30.0, test.s('world'), expires=10)

    # Executes every Monday morning at 7:30 a.m.
    sender.add_periodic_task(
        crontab(hour=7, minute=30, day_of_week=1),
        test.s('Happy Mondays!'),
    )

@app.task
def test(arg):
    print(arg)

问题

但是如果我们使用 django 并且我们的任务被放置在另一个应用程序中呢?使用 celery 4.0.1 我们不再有 @periodic_task 装饰器。所以让我们看看我们能做些什么。

第一种情况

如果您希望将任务及其日程安排彼此靠近:

main_app/some_app/tasks.py

from main_app.celery import app as celery_app

@celery_app.on_after_configure.connect
    def setup_periodic_tasks(sender, **kwargs):
        # Calls test('hello') every 10 seconds.
        sender.add_periodic_task(10.0, test.s('hello'))

@celery_app.task
def test(arg):
    print(arg)

我们可以在调试模式下运行beat

celery -A main_app beat -l debug

我们会看到没有这样的周期性任务。

第二种情况

我们可以尝试在配置文件中这样描述所有周期性任务:

main_app/celery.py

...
app = Celery()

@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    # Calls test('hello') every 10 seconds.
    from main_app.some_app.tasks import test
    sender.add_periodic_task(10.0, test.s('hello'))
...

结果是一样的。但是它的行为会有所不同,您可以通过 pdb 手动调试看到。在第一个示例中,setup_periodic_tasks 回调根本不会被触发。但在第二个例子中,我们会得到django.core.exceptions.AppRegistryNotReady: Apps aren't loaded yet.(这个异常不会被打印出来)

【问题讨论】:

  • SO 绝对欢迎以问答形式分享信息。但是,您在这里得到的并不是一个写得很好的问题。请重写此代码,使其读作面临实际问题的人所写的实际问题。您已经知道解决方案,但是从不知道的人的角度写问题。 (在我看来,您可以从 p.o.v. 中提出一个问题,即某人从 3.x 迁移到 4.x 并发现过去有效的方法不再有效。)
  • 此外,从“问题”标题开始的所有内容都是一个解决方案,应该是一个正式的答案。 (您可以发布截然不同的解决方案作为不同的答案。然后人们可以独立投票。)

标签: python django celery


【解决方案1】:

对于 django,我们需要使用另一个信号:@celery_app.on_after_finalize.connect。两者都可以使用:

  • app/tasks.py 中声明接近任务的任务计划,因为在导入所有tasks.py 并且所有可能的接收者都已订阅(第一种情况)之后,将触发此信号。
  • 集中调度声明,因为 django 应用程序已经初始化并准备好导入(第二种情况)

我想我应该写下最终声明:

第一种情况

任务计划接近任务的声明:

main_app/some_app/tasks.py

from main_app.celery import app as celery_app

@celery_app.on_after_finalize.connect
    def setup_periodic_tasks(sender, **kwargs):
        # Calls test('hello') every 10 seconds.
        sender.add_periodic_task(10.0, test.s('hello'))

@celery_app.task
def test(arg):
    print(arg)

第二种情况

配置文件main_app/celery.py中的集中调度声明:

...

app = Celery()

@app.on_after_finalize.connect
def setup_periodic_tasks(sender, **kwargs):
    # Calls test('hello') every 10 seconds.
    from main_app.some_app.tasks import test
    sender.add_periodic_task(10.0, test.s('hello'))
...

【讨论】:

  • 这仍然无法与人们声称的 @shared_task 一起使用...github.com/celery/celery/issues/3797
  • @HemanthSP,检查docs.celeryproject.org/en/latest/userguide/… 对于这个例子,你可以使用celery -A main_app beat -l debug 来运行调度器并运行一个worker celery -A main_app worker -l debug
  • 我的问题是在 other_app/tasks.py 中使用app = Celery()。使用from main_app.celery import app as celery_app 解决了这个问题!
【解决方案2】:

如果意图是在tasks.py中单独维护任务逻辑,那么在setup_periodic_tasks中调用from main_app.some_app.tasks import test对我不起作用。有效的方法如下:

celery.py

@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    # Calls test('hello') every 10 seconds.
    sender.add_periodic_task(10.0, test.s('hello'), name='add every 10')

@app.task
def test(arg):
    print(arg)
    from some_app.tasks import test
    test(arg)

tasks.py

@shared_task
def test(arg):
    print('world')

这导致以下输出:

[2017-10-26 22:52:42,262: INFO/MainProcess] celery@ubuntu-xenial ready.
[2017-10-26 22:52:42,263: INFO/MainProcess] Received task: main_app.celery.test[3cbdf4fa-ff63-401a-a9e4-cfd1b6bb4ad4]  
[2017-10-26 22:52:42,367: WARNING/ForkPoolWorker-2] hello
[2017-10-26 22:52:42,368: WARNING/ForkPoolWorker-2] world
[2017-10-26 22:52:42,369: INFO/ForkPoolWorker-2] Task main_app.celery.test[3cbdf4fa-ff63-401a-a9e4-cfd1b6bb4ad4] succeeded in 0.002823335991706699s: None
[2017-10-26 22:52:51,205: INFO/Beat] Scheduler: Sending due task add every 10 (main_app.celery.test)
[2017-10-26 22:52:51,207: INFO/MainProcess] Received task: main_app.celery.test[ce0f3cfc-54d5-4d74-94eb-7ced2e5a6c4b]  
[2017-10-26 22:52:51,209: WARNING/ForkPoolWorker-2] hello
[2017-10-26 22:52:51,209: WARNING/ForkPoolWorker-2] world

【讨论】:

  • 这是唯一对我有用的东西。谢谢。
【解决方案3】:

如果您想单独使用任务逻辑,请使用此设置:

celery.py

import os
from celery import Celery
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'backend.settings') # your settings.py path

app = Celery()

@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    sender.add_periodic_task(5, periodic_task.s('sms'), name='SMS Process')
    sender.add_periodic_task(60, periodic_task.s('email'), name='Email Process')


@app.task
def periodic_task(taskname):
    from myapp.tasks import sms_process, email_process

    if taskname == 'sms':
        sms_process()

    elif taskname == 'email':
        email_process()

一个名为 myapp 的 django 应用中的示例任务:

myapp/tasks.py

def sms_process():
    print('send sms task')

def email_process():
    print('send email task')

【讨论】:

    【解决方案4】:

    我使用它来工作

    芹菜.py

    import os
    from celery import Celery
    
    os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'mysite.settings')
    
    app = Celery('mysite')
    app.config_from_object('django.conf:settings', namespace='CELERY')
    app.autodiscover_tasks()
    

    tasks.py

    from celery import current_app
    app = current_app._get_current_object()
    
    @app.task
    def test(arg):
        print(arg)
    
    @app.on_after_finalize.connect
    def app_ready(**kwargs):
        """
        Called once after app has been finalized.
        """
        sender = kwargs.get('sender')
    
        # periodic tasks
        speed = 5
        sender.add_periodic_task(speed, test.s('foo'),name='update leases every {} seconds'.format(speed))
    

    以工人身份运行

    celery -A mysite worker --beat --scheduler django --loglevel=info
    

    【讨论】:

      【解决方案5】:

      也在苦苦挣扎,终端没有任何活动,可以在下面使用:

      Django 3.2.8 版,Celery 5.2.0 版

      在 Django 项目中,称为Proj

      Proj/Proj celery.pysettings.py 旁边的文件)

      celery.py

      import os
      
      from celery import Celery
      
      # Set the default Django settings module for the 'celery' program.
      os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'Proj.settings')
      
      app = Celery('Proj')
      
      # Using a string here means the worker doesn't have to serialize
      # the configuration object to child processes.
      # - namespace='CELERY' means all celery-related configuration keys
      #   should have a `CELERY_` prefix.
      app.config_from_object('django.conf:settings', namespace='CELERY')
      
      # Load task modules from all registered Django apps.
      app.autodiscover_tasks()
      

      __init__.py 内(与settings.py 相同的文件夹)

      __init__.py

      # 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
      
      __all__ = ('celery_app',)
      

      在任何子 django 应用文件夹中,有一个名为 tasks.py 的文件(models.py 旁边)

      tasks.py

      from Proj.celery import app
      
      # Schedule
      @app.on_after_finalize.connect
      def setup_periodic_tasks(sender, **kwargs):
          # Calls test('hello') every 1 seconds.
          sender.add_periodic_task(1.0, test.s('hello'), name='add every 1')
      
          # Calls test('world') every 3 seconds
          sender.add_periodic_task(3.0, test.s('world'), expires=10)
      
      # Tasks
      @app.task
      def test(arg):
          print(arg)
      

      然后,在终端中运行以下命令,使用虚拟环境(如果适用):

      >>> celery -A Proj worker -B
      

      RESULT(确认正常):

      [2021-11-10 11:22:22,070: WARNING/MainProcess] /.venv/lib/python3.9/site-packages/celery/fixups/django.py:203: UserWarning: Using settings.DEBUG leads to a memory
                  leak, never use this setting in production environments!
        warnings.warn('''Using settings.DEBUG leads to a memory
      
      [2021-11-10 11:22:22,173: WARNING/ForkPoolWorker-9] hello
      [2021-11-10 11:22:22,173: WARNING/ForkPoolWorker-3] hello
      [2021-11-10 11:22:22,173: WARNING/ForkPoolWorker-2] world
      

      【讨论】:

        猜你喜欢
        • 2018-11-28
        • 1970-01-01
        • 2020-09-30
        • 1970-01-01
        • 2013-12-05
        • 2019-02-16
        • 2012-01-03
        • 2016-02-04
        • 1970-01-01
        相关资源
        最近更新 更多