Scrapy:如何在Pipeline中将item传递给close_spider方法
问题场景
- 爬虫运行时通过
yield产出大量item,Pipeline中已实现逻辑:基于item的字段(包含country国家字段)更新跟踪表 - 需求:爬虫正常结束后,按国家维度拆分跟踪表,发送给对应负责人员
- 当前通过
spider_closed信号捕获爬虫关闭时机,代码如下:
@classmethod def from_crawler(cls, crawler): temp = cls() crawler.signals.connect(temp.customize_close_spider, signal=signals.spider_closed) return temp def customize_close_spider(self, **kwargs): reason = kwargs.get("reason") spider = kwargs.get("spider") if reason == "finished": # 待执行分国别发送逻辑
- 现存问题:
from_crawler和customize_close_spider方法都无法直接接收到item参数,拿不到country维度数据,无法完成分发表单的逻辑。
实现方案
不需要额外修改信号传参,Scrapy的Pipeline在整个爬虫生命周期内是单例实例,直接在Pipeline实例上维护状态缓存即可,步骤如下:
- 初始化实例缓存
在Pipeline的__init__方法中定义容器,缓存爬取过程中收集到的国家维度数据,不需要存全量item,只存后续分发需要的必要字段即可,控制内存占用:
from collections import defaultdict def __init__(self): # 按国家分组存储需要发送的跟踪表条目 self.country_track_data = defaultdict(list) # 增加执行标记,避免信号重复触发时重复发送 self.send_task_done = False
- 处理item时同步更新缓存
在原有process_item方法的跟踪表更新逻辑后,同步提取当前item的国家信息和对应跟踪表字段,存入实例缓存:
def process_item(self, item, spider): # 原有业务逻辑:更新跟踪表 # self.update_track_table(item) # 同步收集维度数据 country = item.get("country") if country: # 仅存入发送通知需要的字段,不要冗余存全量item track_record = { # 从item中提取跟踪表需要发送的字段,比如任务id、爬取数量、更新时间等 } self.country_track_data[country].append(track_record) return item
如果你的跟踪表是落地在数据库/外部存储的,这里甚至不需要存跟踪表条目,只需要用一个
set存爬取过程中出现过的所有country值即可,爬虫关闭时直接按country去数据库查对应跟踪表数据,内存占用更低。
- 关闭回调中直接读取缓存执行发送
customize_close_spider是绑定在当前Pipeline实例上的方法,可以直接访问实例属性,不需要额外传参:
def customize_close_spider(self, **kwargs): reason = kwargs.get("reason") spider = kwargs.get("spider") if reason == "finished" and not self.send_task_done: self.send_task_done = True # 遍历按国家分组好的数据,匹配对应负责人发送即可 for country, records in self.country_track_data.items(): # 实现发送逻辑:匹配country对应负责人,将records整理为报表发送 pass
注意事项
- 不要尝试修改Scrapy内置信号的传参逻辑,内置信号触发时传递的参数是固定的,强行修改容易引发版本兼容问题
- 如果爬取数据量极大,不要在内存中缓存全量跟踪表记录,优先采用「缓存国家维度标识+关闭时查库」的方案,避免内存溢出
- 如果爬虫异常中断(reason不是
finished),可以根据业务需求决定是否执行发送逻辑
内容的提问来源于stack exchange,提问作者Evgeniy_D
相关产品推荐
相关产品推荐

