【问题标题】:How to spawn a celery task from previous celery task?如何从以前的芹菜任务中产生芹菜任务?
【发布时间】:2017-07-27 14:21:49
【问题描述】:

我可能不正确地使用芹菜。但是我正在开发的聊天机器人需要带有 redis 的 celery 来执行异步任务。这是我正在使用的框架:http://microsoftbotframework.readthedocs.io/en/latest/asynctasks/

我的特定用例目前要求我永远运行 celery 任务,并在其间等待一段任意时间,范围从 30 分钟到 3 天。像这样的

@celery.task
def myAsyncMethod():
    while true:
        timeToWait = getTimeToNextAlarm()
        sleep(timeToWait)
        sendOutMessages()

基本上,我有一个永远不会退出的异步进程。我很确定不应该像这样使用芹菜。 所以我的问题是,我如何创建一个 celery 任务来处理第一个任务,生成一个任务并将其提交到 celery 队列并退出。基本上是这样的:

@celery.task
def myImprovedTask():
    timeToWait = getTimeToNextAlarm()
    sleep(timeToWait)
    sendOutMessages()
    myImprovedTask().delay()    # recursive call to async method for next event

不一定是递归的,甚至不是这样的,而是 celery 原本打算使用的方式(我相信用于短期任务?)

Tl;dr:我如何从另一个任务中创建一个 celery 任务并使原始任务退出?

请告诉我是否应该进一步解释。谢谢。

【问题讨论】:

标签: python asynchronous redis celery botframework


【解决方案1】:

如果您想从初始任务运行另一个任务,只需像通常使用 Task.delay()Task.apply_async() 一样调用它:

@celery.task
def myImprovedTask():
    timeToWait = getTimeToNextAlarm()
    sleep(timeToWait)
    sendOutMessages()
    myImprovedTask.delay()

你是否再次调用相同的任务并不重要。它与delay() 一起排队,您的原始任务返回,然后队列中的下一个任务开始运行。


所有这一切都是基于一个假设你实际上是在异步调用你的 Celery 任务。有时情况并非如此,常见的罪魁祸首是task-always-eager 配置选项。默认情况下它是禁用的,但是(来自docs):

如果task_always_eagerTrue所有的任务都会通过阻塞在本地执行,直到任务返回apply_async()Task.delay() 将返回一个 EagerResult 实例,它模拟 AsyncResult 的 API 和行为,但结果已经被评估。

也就是说,任务将在本地执行,而不是发送到队列中。

所以,请确保您的 Celery 配置包括:

task_always_eager = False

【讨论】:

  • 我确实试过这个。其实我应该提到的。我使用了 delay() 本身。但过了一段时间,它的工作人员用完了,没有任何任务正在执行。
  • 您确定您的任务是异步的吗?排队后,delay() 应该会返回,上一个任务就会完成。
  • 让我快速尝试一下,我们会尽快回复您。
  • 如果你的任务是"eager",它们会在本地运行并阻塞直到完成。
  • 我没有明确地将它们设置为急切设置,但似乎是这样。因为,当我执行 `myImprovedTask.delay() logging.info('ending first event async task.bye..')` 时。在我的 pycharm 调试器中,在方法的开头放置了一个调试点。在调用 myImprovedTask.delay() 时,它直接进入该方法内部到该调试点。日志消息未打印在控制台上。
【解决方案2】:

如果任务没有在当前进程中注册,你可以使用 send_task() 改为按名称调用任务

这里的文档中定义http://docs.celeryproject.org/en/latest/reference/celery.html#celery.Celery.send_task

app.send_task('task_name')

这样做,您必须明确命名任务,例如:

@celery.task(name="myImprovedTask")
def myImprovedTask():

然后你可以调用它:

app.send_task('myImprovedTask')

如果你不喜欢这种方式(或者你有文件在同一个文件中),你也可以用apply_asyncdelay这样调用它:

myImprovedTask.delay()
myImprovedTask.apply_async()

【讨论】:

  • 您好,感谢您的建议。但是在调用 myImprovedTask.delay() 时,它直接进入下一个任务而不退出第一个任务。
猜你喜欢
  • 1970-01-01
  • 2014-12-04
  • 2011-12-02
  • 2012-08-11
  • 2014-07-14
  • 2012-12-01
  • 2014-04-16
  • 2018-07-04
  • 2012-08-30
相关资源
最近更新 更多