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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 23:18:47