【问题标题】:Set delay between tasks in group in Celery在 Celery 中设置组内任务之间的延迟
【发布时间】:2015-05-29 15:15:10
【问题描述】:

我有一个 python 应用程序,用户可以在其中启动某个任务。

任务的全部目的是执行给定数量的 POST/GET 请求,并以特定的时间间隔发送到给定的 URL。

所以用户给出 N - 请求数,V - 每秒请求数。

考虑到由于 I/O 延迟,实际 r/s 速度可能更大或更小,设计这样的任务有什么更好的方法。

首先我决定将 Celery 与 Eventlet 一起使用,否则我将需要几十个工作,这是不可接受的。

我的幼稚做法:

  • 客户端使用 task.delay() 启动任务
  • 内部任务我做这样的事情:

    @task
    def task(number_of_requests, time_period):
       for _ in range(number_of_requests):
           start = time.time()
           params_for_concrete_subtask = ...
           # .... do some IO with monkey_patched eventlet requests library
           elapsed = (time.time() - start)
           # If we completed this subtask to fast
           if elapsed < time_period / number_of_requests:
               eventlet.sleep(time_period / number_of_requests)
    

一个工作示例是here

如果我们太快,我们会尝试等待以保持所需的速度。如果我们太慢,从客户的角度来看是可以的。我们不违反请求/第二个要求。但是如果我重新启动 Celery,这会正确恢复吗?

我认为这应该可行,但我认为有更好的方法。 在 Celery 中,我可以定义一个具有特定速率限制的任务,这几乎符合我的需求保证。所以我可以使用 Celery group 功能并写:

@task(rate_limit=...)
def task(...):
    #

task_executor = task.s(number_of_requests, time_period)
group(task_executor(params_for_concrete_task) for params_for_concrete_task in ...).delay()

但在这里我对动态的 rate_limit 进行了硬编码,我看不到改变它的方法。我看到了一个例子:

  task.s(....).set(... params ...)

但是我尝试将rate_limit 传递给set 方法它不起作用。

另一个可能更好的想法是使用 Celery 的周期性任务调度程序。默认执行周期和定期执行的任务是固定的。

我需要能够动态创建任务,这些任务会以特定的速率限制定期运行给定次数。也许我需要运行我自己的调度程序,它将从数据库中获取任务?但我没有看到任何关于此的文档。

另一种方法是尝试使用chain 函数,但我无法弄清楚任务参数之间是否存在延迟。

【问题讨论】:

  • 使用链会等待上一个工作完成,这是你想要的吗?
  • @Maresh 它们将按顺序运行,它们之间没有延迟。延迟是动态的。这取决于任务执行的时间以及客户端指定的最大速度。因此,如果一个任务运行需要 1 秒,而我想要的速度是 0.5 个任务/秒,我就不能立即运行第二个任务。第二个任务需要等待 1 秒。另一方面,如果运行需要 2 秒,那么下一个任务可以立即运行。我可以从一个任务返回freeze_time 并传递给下一个任务,以便在需要时冻结。但我不喜欢那个解决方案。

标签: python celery


【解决方案1】:

如果您想动态调整 rate_limit,可以使用以下代码进行。它还在运行时创建了chain()。 运行这个你会看到我们成功地将rate_limit 5/sec 改写为0.5/sec。

test_tasks.py

from celery import Celery, signature, chain
import datetime as dt

app = Celery('test_tasks')
app.config_from_object('celery_config')

@app.task(bind=True, rate_limit=5)
def test_1(self):
    print dt.datetime.now()


app.control.broadcast('rate_limit',
                       arguments={'task_name': 'test_tasks.test_1',
                                  'rate_limit': 0.5})

test_task = signature('test_tasks.test_1').set(immutable=True)

l = [test_task] * 100

chain = chain(*l)
res = chain()

我也尝试从类中覆盖该属性,但 IMO 是在工作人员注册任务时设置 rate_limit,这就是 .set() 无效的原因。我在这里推测,必须检查源代码。

解决方案 2

使用上一个调用的结束时间实现自己的等待机制,在链中函数的返回被传递给下一个。

所以它看起来像这样:

from celery import Celery, signature, chain
import datetime as dt
import time

app = Celery('test_tasks')
app.config_from_object('celery_config')

@app.task(bind=True)
def test_1(self, prev_endtime=dt.datetime.now(), wait_seconds=5):
    wait = dt.timedelta(seconds=wait_seconds)
    print dt.datetime.now() - prev_endtime
    wait = wait - (dt.datetime.now() - prev_endtime)
    wait = wait.seconds
    print wait
    time.sleep(max(0, wait))
    now = dt.datetime.now()
    print now
    return now

#app.control.rate_limit('test_tasks.test_1', '0.5')
test_task = signature('test_tasks.test_1')

l = [test_task] * 100

chain = chain(*l)
res = chain()

我认为这实际上比广播更可靠。

【讨论】:

  • 谢谢。我会试试这个。但它应该是“test_tasks.test_1”而不是test_tasks.mytask
  • 是的,很好 :) 我只是想到了另一种可能更好的方法。今天晚些时候我会编辑
  • 查看其他解决方案。
  • 是的,下一个解决方案看起来像我实现的东西。所以我只是创建一个不同的签名,比如test_task = test_tasks.test_1.s(wait_seconds=10)?它会在第一次运行之间增加一个延迟,但这还不错。
  • 是的。如果您想阻止第一次运行等待,只需将旧的 prev_endtime(现在 - X 秒)传递给函数。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-09-28
  • 2011-03-22
  • 1970-01-01
  • 1970-01-01
  • 2014-07-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多