【问题标题】:Run a Scrapy spider in a Celery Task在 Celery 任务中运行 Scrapy 蜘蛛
【发布时间】:2014-04-02 17:08:42
【问题描述】:

This is not working anymore,scrapy 的 API 变了。

现在文档提供了“Run Scrapy from a script”的方法,但我收到了ReactorNotRestartable 错误。

我的任务:

from celery import Task

from twisted.internet import reactor

from scrapy.crawler import Crawler
from scrapy import log, signals
from scrapy.utils.project import get_project_settings

from .spiders import MySpider



class MyTask(Task):
    def run(self, *args, **kwargs):
        spider = MySpider
        settings = get_project_settings()
        crawler = Crawler(settings)
        crawler.signals.connect(reactor.stop, signal=signals.spider_closed)
        crawler.configure()
        crawler.crawl(spider)
        crawler.start()

        log.start()
        reactor.run()

【问题讨论】:

  • 你用的是什么版本的scrapy?
  • @Talvalin Scrapy==0.22.2
  • @shirkey 我在第一个链接中提到了那个问题

标签: scrapy twisted celery


【解决方案1】:

扭曲的反应堆无法重新启动。解决此问题的方法是让 celery 任务为您要执行的每个爬网创建一个新的子进程,如以下帖子中所建议的那样:

这通过使用multiprocessing 包解决了“reactor 无法重新启动”问题。但问题在于,由于您将遇到另一个问题,即守护进程无法生成子进程,因此最新的 celery 版本现在已经过时了。因此,为了使解决方法起作用,您需要使用 celery 版本。

是的,scrapy API 已更改。但稍作修改(import Crawler 而不是CrawlerProcess)。您可以通过使用 celery 版本来获得解决方法。

可以在此处找到 Celery 问题: Celery Issue #1709

这是我的更新的爬虫脚本,通过使用billiard 而不是multiprocessing,与较新的芹菜版本一起工作:

from scrapy.crawler import Crawler
from scrapy.conf import settings
from myspider import MySpider
from scrapy import log, project
from twisted.internet import reactor
from billiard import Process
from scrapy.utils.project import get_project_settings
from scrapy import signals


class UrlCrawlerScript(Process):
    def __init__(self, spider):
        Process.__init__(self)
        settings = get_project_settings()
        self.crawler = Crawler(settings)
        self.crawler.configure()
        self.crawler.signals.connect(reactor.stop, signal=signals.spider_closed)
        self.spider = spider

    def run(self):
        self.crawler.crawl(self.spider)
        self.crawler.start()
        reactor.run()

def run_spider(url):
    spider = MySpider(url)
    crawler = UrlCrawlerScript(spider)
    crawler.start()
    crawler.join()

编辑:通过阅读 celery 问题#1709,他们建议使用台球而不是多处理,以便解除子进程限制。换句话说,我们应该尝试billiard 看看它是否有效!

编辑 2: 是的,通过使用 billiard,我的脚本适用于最新的 celery 版本!查看我更新的脚本。

【讨论】:

  • 注意 - 我必须将 self.crawler.signals.connect(reactor.stop, signal=signals.spider_closed) 行移到初始化检查之外,否则第二次运行会挂起。移动它可以使它在我的项目中正常工作。此外,随着scrapy.project 的贬值,使用台球的current_thread 在每个线程的基础上设置初始化标志。这也很有效。
  • jlovison,您能分享一下您在 current_thread 上所做的更改吗?你在哪里放置了signals.spider_closed?提前致谢
  • @BjBlazkowicz 因为 Process 是这里的基类,并且调用了 Process.__init__(self) ,派生类 UrlCrawlerScript 不是 also necessary for __del__ to be called 还是会自动调用它?
  • @lennard 如果不分叉一个新进程,它将无法工作。它的工作时间是你设置的并发的 x 倍,但每个 celery worker 只能工作一次。但是,如果您可以在每个任务之后加入 celery 进程,它就会起作用。为此,您可以使用以下设置:CELERYD_MAX_TASKS_PER_CHILD = 1
  • 对于celery==4.1.0 Scrapy==1.5.0 billiard==3.5.0.3,我尝试对此进行修改但失败了。我在 django 中使用它。然后我尝试了CrawlerRunner,也失败了。最终我放弃了,转而使用CELERY_WORKER_MAX_TASKS_PER_CHILD = 1。代码在this gist
【解决方案2】:

Twisted reactor 无法重新启动,因此一旦一个蜘蛛程序完成运行并且crawler 隐式停止了 reactor,该 worker 将毫无用处。

正如在另一个问题的答案中所发布的那样,您需要做的就是杀死运行您的蜘蛛的工人并用一个新的替换它,这样可以防止反应堆多次启动和停止。为此,只需设置:

CELERYD_MAX_TASKS_PER_CHILD = 1

缺点是您并没有真正使用 Twisted reactor 发挥其全部潜力并浪费运行多个反应器的资源,因为一个反应器可以在一个进程中同时运行多个蜘蛛。更好的方法是每个工人运行一个反应器(甚至全球一个反应器)并且不要让crawler 触摸它。

我正在为一个非常相似的项目做这个,所以如果我有任何进展,我会更新这篇文章。

【讨论】:

  • 我对您的解决方法很感兴趣。当您有什么想法时,请告诉我们。
  • 非常有趣,你用这个有什么收获吗?
【解决方案3】:

为了避免在 Celery 任务队列中运行 Scrapy 时出现 ReactorNotRestartable 错误,我使用了线程。在一个应用程序中多次运行 Twisted reactor 的方法相同。 Scrapy 也使用了 Twisted,所以我们可以这样做。

代码如下:

from threading import Thread
from scrapy.crawler import CrawlerProcess
import scrapy

class MySpider(scrapy.Spider):
    name = 'my_spider'


class MyCrawler:

    spider_settings = {}

    def run_crawler(self):

        process = CrawlerProcess(self.spider_settings)
        process.crawl(MySpider)
        Thread(target=process.start).start()

不要忘记为 celery 增加 CELERYD_CONCURRENCY。

CELERYD_CONCURRENCY = 10

对我来说很好。

这不是阻塞进程运行,但无论如何,scrapy 的最佳实践是在回调中处理数据。这样做:

for crawler in process.crawlers:
    crawler.spider.save_result_callback = some_callback
    crawler.spider.save_result_callback_params = some_callback_params

Thread(target=process.start).start()

【讨论】:

    【解决方案4】:

    如果您有很多任务要处理,我会说这种方法非常低效。 因为 Celery 是线程化的 - 在它自己的线程中运行每个任务。 假设使用 RabbitMQ 作为代理,您可以通过 >10K q/s。 使用 Celery,这可能会导致 10K 线程开销! 我建议不要在这里使用芹菜。而是直接访问代理!

    【讨论】:

    • 直接访问代理?什么意思?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-02-15
    • 1970-01-01
    相关资源
    最近更新 更多