【问题标题】:how to chain celery tasks如何链接芹菜任务
【发布时间】:2015-02-05 06:50:22
【问题描述】:

我想以标准方式链接 celery 任务。

我有一个 json 文件。在该文件中,有许多硬编码的 url。我需要废弃这些链接,并废弃在抓取这些链接时发现的链接。

目前,我正在这样做。

for each_news_source, news_categories in rss_obj.iteritems():
    for each_category in news_categories:
        category = each_category['category']
        rss_link = each_category['feed']
        json_id = each_category['json']
        try:
            list_of_links = getrsslinks(rss_link)
            for link in list_of_links:
                scrape_link.delay(link, json_id, category)
        except Exception,e:
            print "Invalid url", str(e)

我想要 getrsslinks 也是 celery 任务的东西,然后报废 getrsslinks 返回的 url 列表也应该是另一个 celery 任务。

它遵循这种模式

harcodeJSONURL1--
               --`getrsslinks` (celery task)
                               --scrap link 1 (celery task)
                               --scrap link 2 (celery task)
                               --scrap link 3 (celery task)
                               --scrap link 4 (celery task)

harcodeJSONURL2--
               --`getrsslinks` (celery task)
                               --scrap link 1 (celery task)
                               --scrap link 2 (celery task)
                               --scrap link 3 (celery task)
                               --scrap link 4 (celery task)

等等..

我该怎么做?

【问题讨论】:

    标签: python celery


    【解决方案1】:

    看看 Celery 中的 subtask options。在您的案例组中应该有所帮助。您只需要在getrsslinks 中调用scrape_link 组。

    from celery import group
    
    @app.task
    def getrsslinks(rsslink, json_id, category):
        # do processing
    
        # Call scrape links
        scrape_jobs = group(scrape_link.s(link, json_id, category) for link in link_list)
        scrape_jobs.apply_async()
        ...
    

    您可能希望getrsslinks 返回scrape_jobs 以更轻松地监控作业。然后在解析你的json文件时,你会像这样调用getrsslinks

    for each_news_source, news_categories in rss_obj.iteritems():
        for each_category in news_categories:
            category = each_category['category']
            rss_link = each_category['feed']
            json_id = each_category['json']
            getrsslinks.delay(rss_link, json_id, category)
    

    最后,要监控哪些链接无效(因为我们替换了 try/except 块),您需要存储所有 getrsslinks 任务并观察成功或失败。为此,您可以使用 apply_asynclink_error

    【讨论】:

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