【问题标题】:Celery task calling itself after task succeeds without celerybeatCelery 任务在没有 celerybeat 的情况下成功后调用自身
【发布时间】:2019-09-08 09:58:55
【问题描述】:

我想在当前任务完成后每隔 30 分钟调用一次我的 celery 任务,但有时任务需要的时间比预期的要长,因为任务是从远程服务器下载文件。所以我不想使用 celeryBeat。另外,使用自我。重试仅适用于我想发生错误时。这是我的任务:

@shared_task(name="download_big", bind=True, acks_late=true, autoretry_for=(Exception, requests.exceptiosn.RequestException), retry_kwargs={"max_retries": 4, "countdown": 3}):
def download_big(self):
    my_file = session.get('example.com/hello.mp4')
    if my_file.status_code == requests.codes["OK"]:
        open("hello.mp4", "wb").write(my_file.content)
    else:
        self.retry()

更新:

好吧,我将结构更改为:

@shared_task(name="download_big", bind=True, acks_late=true, autoretry_for=(Exception, requests.exceptiosn.RequestException), retry_kwargs={"max_retries": 4, "countdown": 3}):
def download_big(url):
    my_file = session.get(url, name)
    if my_file.status_code == requests.codes["OK"]:
        open(name, "wb").write(my_file.content)
    else:
        self.retry()

@shared_task(name="download_all", bind=True, acks_late=true, autoretry_for=(Exception, requests.exceptiosn.RequestException), retry_kwargs={"max_retries": 4, "countdown": 3}):
def download_all(self):
    my_list = [...]  # bunch of urls with names
    jobs = []
    for name, url in my_list:
        jobs.append(download_big.si(url, name))
    group(jobs)()

所以在这种情况下,我必须调用 download_all 方法而不是 download_big,这样我可以并行下载文件,并且当所有组任务完成后,它需要在 30 分钟后再次调用自身。

【问题讨论】:

    标签: python django celery django-celery


    【解决方案1】:

    您可以尝试使用chord,它将运行一组任务,当它们完成时,将运行一个回调,您可以使用它来重新安排。

    例如

    from celery import chord
    
    @shared_task(name="download_big", bind=True, acks_late=true, autoretry_for=(Exception, requests.exceptiosn.RequestException), retry_kwargs={"max_retries": 4, "countdown": 3}):
    def download_big(url):
        my_file = session.get(url, name)
        if my_file.status_code == requests.codes["OK"]:
            open(name, "wb").write(my_file.content)
        else:
            self.retry()
    
    @shared_task(name="download_all", bind=True, acks_late=true, autoretry_for=(Exception, requests.exceptiosn.RequestException), retry_kwargs={"max_retries": 4, "countdown": 3}):
    def download_all(self):
        my_list = [...]  # bunch of urls with names
        jobs = []
        for name, url in my_list:
            jobs.append(download_big.si(url, name))
    
        # Run the group and reschedule once all tasks complete
        chord(jobs)(download_all.apply_async(countdown=1800))
    

    【讨论】:

    • 呃,下载没完成会自己调用吗?如果不是在此任务中下载文件,而是调用下载任务组,情况如何?那样组调用然后文件都没有完成,download_big 再次调用自己?
    • 好吧,调度本身的调用发生在文件写入之后,我认为这意味着下载完成?对于download_all(),您可以使用和弦,它允许您在任务完成时指定回调。在回调中,您可以重新安排组。我会更新答案。
    • 这把Retry in 4s: AttributeError("'AsyncResult' object has no attribute 'clone'",扔给我
    猜你喜欢
    • 2012-12-26
    • 2022-01-21
    • 1970-01-01
    • 2023-03-10
    • 2017-01-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-05-12
    相关资源
    最近更新 更多