如何在Scrapy爬虫中异步处理DataFrame的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

