Scrapy使用SQLAlchemy1.4 ORM实现数据库批量插入的正确方法
SQLAlchemy 1.4 + Scrapy 批量入库实现方案
原有方案问题点
- 单条逐次插入时每条数据都单独触发事务、查询、提交、连接关闭,IO和事务开销占比超过90%,入库效率极低
- models层直接初始化全局session,在Scrapy并发爬取场景下会出现跨线程事务冲突、连接状态异常,是之前bulk操作报错的核心原因
- 全量缓存所有item到spider关闭才写入,数据量大时会占满内存,爬取中途中断则所有缓存数据全部丢失
- 全量bulk方案未做重复校验,直接插入会触发
Reference字段的重复值报错,传入Scrapy Item对象而非普通字典也会导致bulk_insert_mappings参数校验失败
第一步:修正models层连接配置
去掉全局session实例,只保留连接引擎和会话工厂,同时给去重字段加唯一索引提升查重效率:
from sqlalchemy import Column, String, Integer, create_engine from sqlalchemy.orm import sessionmaker from sqlalchemy.ext.declarative import declarative_base from . import settings engine = create_engine( settings.DATABSE_URL, pool_pre_ping=True, pool_size=10, max_overflow=20 ) SessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False) DeclarativeBase = declarative_base() class Olx_Eg(DeclarativeBase): __tablename__ = "olx_egypt" _id = Column(Integer, primary_key=True) URL = Column("URL", String) Breadcrumb = Column("Breadcrumb", String) Price = Column("Price", String) Title = Column("Title", String) Type = Column("Type", String) Bedrooms = Column("Bedrooms", String) Bathrooms = Column("Bathrooms", String) Area = Column("Area", String) Location = Column("Location", String) Compound = Column("Compound", String) seller = Column("seller", String) Seller_member_since = Column("Seller_member_since", String) Seller_phone_number = Column("Seller_phone_number", String) Description = Column("Description", String) Amenities = Column("Amenities", String) # 加唯一索引,兼顾查重性能和库层面防重 Reference = Column("Reference", String, unique=True) Listed_date = Column("Listed_date", String) Level = Column("Level", String) Payment_option = Column("Payment_option", String) Delivery_term = Column("Delivery_term", String) Furnished = Column("Furnished", String) Delivery_date = Column("Delivery_date", String) Down_payment = Column("Down_payment", String) Image_url = Column("Image_url", String) # 首次运行自动建表 DeclarativeBase.metadata.create_all(bind=engine)
第二步:实现分批插入Pipeline
设置合理批次阈值(推荐100-500条/批),累计到阈值就触发批量写入,兼顾内存占用、写入性能和数据安全性:
from olx_egypt.models import Olx_Eg, SessionLocal class OlxEgBatchPipeline: def __init__(self, batch_size=200): self.batch_size = batch_size self.items_cache = [] self.session = None def open_spider(self, spider): # 每个spider启动时创建独立会话,避免全局session冲突 self.session = SessionLocal() def process_item(self, item, spider): # 转成普通字典,适配bulk_insert_mappings参数要求 self.items_cache.append(dict(item)) if len(self.items_cache) >= self.batch_size: self._flush_batch(spider) return item def _flush_batch(self, spider): if not self.items_cache: return try: # 批次内去重 ref_set = set() unique_items = [] for it in self.items_cache: ref = it.get("Reference") if ref not in ref_set: ref_set.add(ref) unique_items.append(it) # 查询库中已存在的Reference,过滤重复数据 exist_refs = set( r[0] for r in self.session.query(Olx_Eg.Reference) .filter(Olx_Eg.Reference.in_(ref_set)) .all() ) to_insert = [it for it in unique_items if it.get("Reference") not in exist_refs] # 批量写入 if to_insert: self.session.bulk_insert_mappings(Olx_Eg, to_insert) self.session.commit() except Exception as e: self.session.rollback() spider.logger.error(f"批次写入失败: {str(e)}") finally: self.items_cache.clear() def close_spider(self, spider): # 把剩余不足一个批次的缓存数据写入 self._flush_batch(spider) self.session.close()
性能优化说明
- 若使用MySQL数据库,在连接URL末尾追加*
?rewriteBatchedStatements=true*参数,驱动层会自动把多条插入合并成多值INSERT语句,写入性能可再提升3-5倍;PostgreSQL无需额外配置即可享受批量写入优化 - 批次大小按需调整:如果单条数据包含长文本字段,适当调小batch_size到100-200,避免单次提交数据包过大导致数据库超时
- 不要实例化ORM对象再传入bulk方法,直接传字典能省去ORM对象初始化的额外开销,写入速度更快
- 不要跨spider、跨线程复用session,每个Pipeline实例独立维护会话即可,避免事务锁冲突和连接泄漏
- 提前过滤掉item中模型未定义的字段,否则会触发
bulk_insert_mappings字段不存在的报错
内容的提问来源于stack exchange,提问作者dougj
相关产品推荐
相关产品推荐

