如何在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()
关键优化点说明
- 缓冲机制:通过
item_buffer积累Item,达到batch_size才触发插入,大幅减少数据库连接与交互次数 - 批量操作:使用
executemany替代循环execute,单次请求插入多条数据 - 优化重复ID查询:不再全表扫描
products,而是仅查询当前批量中的product_id,查询效率提升显著 - 统一事务提交:批量插入完成后再执行
commit,减少事务提交开销 - 剩余数据处理:爬虫关闭时强制插入缓冲中剩余的Item,避免数据丢失
注意事项
- 批量大小
batch_size可根据数据库性能调整(建议500-2000之间) - 若爬虫运行中出现异常,需确保事务回滚,避免数据不一致
- 可根据需求添加日志记录,替代print语句便于排查问题
内容的提问来源于stack exchange,提问作者cevreyolu
相关产品推荐
相关产品推荐

