You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Scrapy+Celery任务中按顺序间隔2分钟运行爬虫?

如何调整Scrapy+Celery爬虫任务逻辑实现串行执行并间隔2分钟启动

我当前基于Scrapy与Celery搭建了爬虫任务体系,现有signals.py和tasks.py代码如下,请问如何调整任务逻辑,使爬虫能够依次(一个接一个)运行,且相邻爬虫的启动间隔为2分钟?


现有代码

signals.py

@receiver(post_save, sender=ParseCategoryUrl)
def start_parse_from_category_url(sender, created, instance, **kwargs):
    """
    当CategoryUrl模型实例创建时,触发run_spider_to_pars_ads_list任务的信号接收器
    """
    if instance.url:
        group(
            run_spider_to_pars_ads_list.s(url=instance.url, user=instance.user.id),
        ).apply_async()

tasks.py

@app.task(name="run_spider_to_pars_ads_list")
def run_spider_to_pars_ads_list(
    url: str, user: int, pages: Optional[int] = None
) -> None:
    process.crawl(ListitemsSpider, url=url, user=user, pages=pages)
    d = process.join()
    d.addBoth(lambda _: reactor.stop())
    reactor.run()

调整方案

要实现爬虫串行执行且间隔2分钟启动,需从任务调度方式、爬虫运行逻辑两方面修改:

1. 替换并行调度为串行任务链

当前用group会触发并行任务,改成chain可以让任务按顺序执行;同时通过延迟配置实现2分钟间隔。

修改后的signals.py

from celery import chain
from django.db.models.signals import post_save
from django.dispatch import receiver

# 生产环境建议用Redis/数据库存储任务链ID,替代全局变量
current_task_chain = None

@receiver(post_save, sender=ParseCategoryUrl)
def start_parse_from_category_url(sender, created, instance, **kwargs):
    """当CategoryUrl实例创建时,将爬虫任务加入串行执行链,间隔2分钟启动下一个"""
    if not (instance.url and created):
        return
    
    spider_task = run_spider_to_pars_ads_list.s(url=instance.url, user=instance.user.id)
    
    if current_task_chain is None:
        # 第一个任务直接执行
        current_task_chain = spider_task.apply_async()
    else:
        # 后续任务:在上一个任务完成后,延迟2分钟执行当前爬虫
        new_chain = chain(current_task_chain, spider_task.si(countdown=120))
        current_task_chain = new_chain.apply_async()

2. 优化爬虫运行逻辑

原tasks.py中重复启动reactor可能导致异常,改为全局单例CrawlerProcess:

修改后的tasks.py

from scrapy.crawler import CrawlerProcess
from scrapy.utils.project import get_project_settings
from celery import app

# 全局仅初始化一次爬虫进程
process = CrawlerProcess(get_project_settings())

@app.task(name="run_spider_to_pars_ads_list")
def run_spider_to_pars_ads_list(
    url: str, user: int, pages: Optional[int] = None
) -> None:
    process.crawl(ListitemsSpider, url=url, user=user, pages=pages)
    # 启动爬虫并在爬取完成后自动停止
    process.start(stop_after_crawl=True)

3. 生产环境补充配置

  • 单队列串行控制:创建专属串行任务队列,给Celery Worker指定-Q spider_serial_queue并设置并发数--concurrency=1,确保同一时间仅一个任务运行。
  • 任务链持久化:用Redis存储当前任务链ID,避免多Worker环境下全局变量失效,示例:
    import redis
    r = redis.Redis()
    
    # 存储任务链ID
    r.set("current_spider_chain", current_task_chain.id)
    # 获取任务链ID
    prev_chain_id = r.get("current_spider_chain")
    

内容的提问来源于stack exchange,提问作者Kirill

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 21:51:01