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

如何在Scrapy中实现Item缓冲并批量插入MySQL数据库

Scrapy Pipeline 批量插入数据库优化方案

修改思路

  • 新增缓冲列表暂存待插入Item,设定批量阈值(如1000条)
  • 达到阈值时执行批量插入与批量查询,减少数据库交互次数
  • 爬虫关闭时强制处理剩余未插入的Item,避免数据遗漏
  • 优化原有的重复ID判断逻辑,避免全表扫描

优化后的完整代码

import mysql.connector
from mysql.connector import errorcode

class BatchInsertPipeline:
    def __init__(self, **kwargs):
        self.cnx = self.mysql_connect()
        # 初始化缓冲列表与批量阈值
        self.item_buffer = []
        self.batch_size = 1000  # 可根据数据库性能调整
        # 假设self.conf和self.table、self.table2已提前配置

    def open_spider(self, spider):
        print("spider open")

    def process_item(self, item, spider):
        # 将Item转为字典加入缓冲
        self.item_buffer.append(dict(item))
        # 达到批量阈值时执行插入
        if len(self.item_buffer) >= self.batch_size:
            self.batch_insert()
        return item

    def batch_insert(self):
        if not self.item_buffer:
            return
        
        cursor = self.cnx.cursor()
        try:
            # -------------------------- 批量插入主表(table) --------------------------
            insert_table_query = (
                f"INSERT INTO {self.table} "
                "(rowid, date, listing_id, product_id, product_name, price, url) "
                "VALUES (%(rowid)s, %(date)s, %(listing_id)s, %(product_id)s, %(product_name)s, %(price)s, %(url)s)"
            )
            # 批量执行插入
            cursor.executemany(insert_table_query, self.item_buffer)
            print(f"批量插入主表 {cursor.rowcount} 条数据")

            # -------------------------- 批量判断并插入副表(table2) --------------------------
            # 收集所有待判断的product_id
            product_ids = [item['product_id'] for item in self.item_buffer]
            # 批量查询已存在的product_id(避免全表扫描)
            check_exist_query = (
                "SELECT product_id FROM products WHERE product_id IN (%s)"
            ) % ','.join(['%s'] * len(product_ids))
            
            cursor.execute(check_exist_query, product_ids)
            existing_ids = {row[0] for row in cursor.fetchall()}
            
            # 过滤出需要插入副表的Item
            table2_items = [item for item in self.item_buffer if item['product_id'] not in existing_ids]
            if table2_items:
                insert_table2_query = (
                    f"INSERT INTO {self.table2} "
                    "(product_rowid, date, listing_id, product_id, product_name, price, url) "
                    "VALUES (%(rowid)s, %(date)s, %(listing_id)s, %(product_id)s, %(product_name)s, %(price)s, %(url)s)"
                )
                cursor.executemany(insert_table2_query, table2_items)
                print(f"批量插入副表 {cursor.rowcount} 条数据")

            # 统一提交事务(减少提交次数)
            self.cnx.commit()
            # 清空缓冲
            self.item_buffer = []
        except mysql.connector.Error as err:
            print(f"批量插入失败: {err}")
            self.cnx.rollback()
        finally:
            cursor.close()

    def close_spider(self, spider):
        # 处理剩余未插入的Item
        if self.item_buffer:
            self.batch_insert()
        self.mysql_close()
        print("spider closed")

    def mysql_connect(self):
        try:
            return mysql.connector.connect(**self.conf)
        except mysql.connector.Error as err:
            if err.errno == errorcode.ER_ACCESS_DENIED_ERROR:
                print("用户名或密码错误")
            elif err.errno == errorcode.ER_BAD_DB_ERROR:
                print("数据库不存在")
            else:
                print(err)
            raise  # 抛出异常终止爬虫,避免无连接运行

    def mysql_close(self):
        if self.cnx.is_connected():
            self.cnx.close()

关键优化点说明

  1. 缓冲机制:通过item_buffer积累Item,达到batch_size才触发插入,大幅减少数据库连接与交互次数
  2. 批量操作:使用executemany替代循环execute,单次请求插入多条数据
  3. 优化重复ID查询:不再全表扫描products,而是仅查询当前批量中的product_id,查询效率提升显著
  4. 统一事务提交:批量插入完成后再执行commit,减少事务提交开销
  5. 剩余数据处理:爬虫关闭时强制插入缓冲中剩余的Item,避免数据丢失

注意事项

  • 批量大小batch_size可根据数据库性能调整(建议500-2000之间)
  • 若爬虫运行中出现异常,需确保事务回滚,避免数据不一致
  • 可根据需求添加日志记录,替代print语句便于排查问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 16:55:24