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

如何在Scrapy爬虫中异步处理DataFrame的Item生成

优化Scrapy爬虫50万条数据Item生成性能的方案

嘿,我看你这爬虫的核心逻辑没问题,但处理50万条数据时卡在了同步生成Item的阶段——5.5小时确实够头疼的!问题出在process_csvs里的单线程循环完全没利用Scrapy的异步优势,咱们来一步步拆解优化:

一、先解决数据查询的性能瓶颈

你现在每次循环都要对primary_csv_file_3做过滤查询(p_nutrients = primary_csv_file_3[primary_csv_file_3.id == product.supporting_id]),这是O(n)的操作,50万次循环下来开销极大。先把营养数据提前分组映射,把查询变成O(1):

def process_csvs(self):
    # ... 读取CSV的代码保持不变 ...

    # 关键优化:提前按id分组营养数据,生成映射字典
    nutrient_map = {}
    # 按id分组,一次性处理所有营养数据
    for idx, group in primary_csv_file_3.groupby('id'):
        # 把每组转成需要的字典格式,存在map里
        nutrient_map[idx] = [
            dict(
                alias=row.name,
                value=row.amount,
                unit_of_measure=row.units
            ) for row in group.itertuples()
        ]

    # 后续循环直接从map里取数据,不用再过滤
    for product in primary_csv_file_1.itertuples():
        loader = Loader(item=REDACTEDItem())
        # ... 填充基础字段 ...
        
        # 直接通过map获取营养数据,速度提升N倍
        nutrients = nutrient_map.get(product.supporting_id, [])
        loader.add_value('nutrition', json.dumps(nutrients))
        yield loader.load_item()

这一步就能砍掉大部分查询开销,是最立竿见影的优化。

二、利用Scrapy的异步/并行能力处理Item生成

Scrapy基于Twisted异步框架,但你的Item生成是同步阻塞的。我们可以用线程池+异步延迟调用把Item生成任务放到后台线程执行,不阻塞Scrapy的主 reactor:

方案1:单条Item异步生成

把创建Item的逻辑抽成独立函数,用deferToThread放到线程里执行:

from scrapy.utils.defer import deferToThread

def _create_single_item(self, product, nutrient_map):
    """单独抽出生成Item的逻辑,交给线程处理"""
    loader = Loader(item=REDACTEDItem())
    loader.add_value('url', 'REDACTED')
    loader.add_value('category', product.category)
    loader.add_value('upc', product.upc)
    loader.add_value('brand', product.brand)
    loader.add_value('product_name', product.name)
    
    nutrients = nutrient_map.get(product.supporting_id, [])
    loader.add_value('nutrition', json.dumps(nutrients))
    return loader.load_item()

def process_csvs(self):
    # ... 前面的读取、分组逻辑不变 ...

    # 异步处理每条Item生成
    for product in primary_csv_file_1.itertuples():
        # 把创建Item的任务交给线程,Scrapy主进程可以继续处理其他任务
        item = yield deferToThread(self._create_single_item, product, nutrient_map)
        yield item

方案2:批量Chunk并行处理

如果单条异步还不够快,可以把主DataFrame分成多个Chunk,用线程池并行处理每个Chunk,进一步利用多核CPU:

from concurrent.futures import ThreadPoolExecutor
from scrapy.utils.defer import deferToThread

def _process_chunk(self, chunk, nutrient_map):
    """处理一个数据Chunk,生成一组Item"""
    items = []
    for product in chunk.itertuples():
        loader = Loader(item=REDACTEDItem())
        loader.add_value('url', 'REDACTED')
        loader.add_value('category', product.category)
        loader.add_value('upc', product.upc)
        loader.add_value('brand', product.brand)
        loader.add_value('product_name', product.name)
        
        nutrients = nutrient_map.get(product.supporting_id, [])
        loader.add_value('nutrition', json.dumps(nutrients))
        items.append(loader.load_item())
    return items

def process_csvs(self):
    # ... 前面的读取、分组逻辑不变 ...

    # 把主DataFrame分成多个Chunk,比如每个Chunk 1000条
    chunk_size = 1000
    total_rows = len(primary_csv_file_1)
    chunks = [
        primary_csv_file_1[i:i+chunk_size] 
        for i in range(0, total_rows, chunk_size)
    ]

    # 用线程池并行处理Chunk,max_workers根据CPU核心数调整
    with ThreadPoolExecutor(max_workers=4) as executor:
        for chunk in chunks:
            # 提交Chunk处理任务到线程池
            items = yield deferToThread(self._process_chunk, chunk, nutrient_map)
            for item in items:
                yield item

三、额外优化:跳过磁盘IO,直接内存处理CSV

你现在把Zip文件和CSV都写到磁盘再读取,磁盘IO是很大的性能开销。可以直接在内存里处理Zip和CSV,完全跳过磁盘写入:

修改handle_zip方法,用BytesIO内存流处理:

from io import BytesIO

def handle_zip(self, response, supporting_file=None):
    file_alias = 'primary_csv' if supporting_file else 'supporting_csv'
    self.logger.info(f"Processing {file_alias} zip file")
    # 把响应内容读到内存流里,不用保存到磁盘
    zip_stream = BytesIO(response.body)
    
    with zipfile.ZipFile(zip_stream, 'r') as zfile:
        if supporting_file:
            # 直接从Zip流里读取CSV到DataFrame,跳过磁盘写入
            with zfile.open('primary_csv_file_1.csv') as f:
                self.primary_csv_1 = pd.read_csv(f, usecols=[...], dtype=dict(...))
            with zfile.open('primary_csv_file_2.csv') as f:
                self.primary_csv_2 = pd.read_csv(f, usecols=[...], dtype=dict(...))
            with zfile.open('primary_csv_file_3.csv') as f:
                self.primary_csv_3 = pd.read_csv(f, usecols=[...], dtype=dict(...))
        else:
            with zfile.open('supporting_csv_file.csv') as f:
                self.supporting_csv = pd.read_csv(f, usecols=[...], dtype=dict(...))
    
    if supporting_file:
        self.logger.info('Downloading supporting_csv zip file')
        yield scrapy.Request(supporting_file, callback=self.handle_zip)
    else:
        self.logger.info('Processing CSV data')
        yield from self.process_csvs()

然后修改process_csvs,直接用实例变量里的DataFrame,不用再从磁盘读取:

def process_csvs(self):
    primary_csv_file_1 = self.primary_csv_1
    primary_csv_file_2 = self.primary_csv_2
    primary_csv_file_3 = self.primary_csv_3
    supporting_csv_file = self.supporting_csv

    # ... 后续的合并、分组逻辑不变 ...

这一步能节省大量的磁盘读写时间,尤其是在IO性能一般的机器上效果明显。

四、其他小优化

  • 把json.dumps的操作提前批量处理,比如在分组营养数据的时候直接转成字符串,避免循环里重复序列化:
    nutrient_map[idx] = json.dumps([
        dict(
            alias=row.name,
            value=row.amount,
            unit_of_measure=row.units
        ) for row in group.itertuples()
    ])
    
    然后循环里直接loader.add_value('nutrition', nutrients)即可。
  • 如果你的Item Loader逻辑不复杂,可以直接构造Item对象,跳过Loader的开销(不过Loader的优势是数据清洗,看你的需求取舍)。

按照这些优化步骤下来,5.5小时的耗时应该能降到几十分钟甚至更少,完全发挥Scrapy的异步优势!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 11:27:45