批量提交时如何用SQLAlchemy检查会话中待提交数据
批量提交爬虫数据时避免批次内重复的解决方案
问题场景
我自己写了网络爬虫填充数据库,为了减少数据库提交次数采用批量提交方式。目前用以下代码在添加数据前检查是否已存在于表中:
class ModelMixin(Model): __abstract__ = True created_at = Column(DateTime, default=datetime.utcnow()) updated_at = Column(DateTime, onupdate=datetime.utcnow()) @classmethod def get_or_create(cls, session, **kwargs): entry = session.query(cls).filter_by(**kwargs).first() if not entry: logger.info('Inserting') entry = cls(**kwargs) session.add(entry) return entry
遇到的问题
当批次内存在重复数据时,因为数据还没提交到数据库,上面的检查逻辑只会查询数据库,不会识别session中待提交的pending数据,最终提交时触发唯一约束冲突。我需要实现检查预提交批次中是否已存在该数据的逻辑,比如判断if entry not in session.pending。
补充情况
尝试用session.get解决但遇到问题:
>>> data = {'id': 315253, 'html': '<h1>lorem ipsum</h1>'} >>> session.add(DataLake(**data)) >>> >>> session.new IdentitySet([DataLake, ID: 315253]) >>> >>> session.get(DataLake, 315253) # 刚添加的ID返回None >>> session.get(DataLake, 92742) # 数据库中存在的ID返回实例 DataLake, ID: 92742
解决方法
1. 修改get_or_create,先检查session待提交实例
在查询数据库之前,先遍历session中已添加但未提交的实例(session.new集合),根据模型的唯一标识字段判断是否已存在:
class ModelMixin(Model): __abstract__ = True created_at = Column(DateTime, default=datetime.utcnow()) updated_at = Column(DateTime, onupdate=datetime.utcnow()) @classmethod def get_or_create(cls, session, **kwargs): # 替换成你的模型的唯一约束字段,比如'id' unique_field = 'id' unique_value = kwargs.get(unique_field) if unique_value: # 检查session中已添加但未提交的实例 for instance in session.new: if isinstance(instance, cls) and getattr(instance, unique_field) == unique_value: return instance # 再查询数据库 entry = session.query(cls).filter_by(**kwargs).first() if not entry: logger.info('Inserting') entry = cls(**kwargs) session.add(entry) return entry
2. 批次处理时用临时缓存跟踪(高效方案)
如果批次数据量较大,遍历session.new效率偏低,可以在批次处理时维护一个字典缓存已添加的实例,避免重复检查:
# 批次处理时的缓存,键为模型唯一标识值 batch_cache = {} def process_crawler_data(session, data_list): for data in data_list: # 替换成你的模型类 model_cls = DataLake unique_id = data['id'] # 先查缓存,存在则跳过 if unique_id in batch_cache: continue # 调用get_or_create添加数据 entry = model_cls.get_or_create(session, **data) batch_cache[unique_id] = entry # 批量提交 session.commit()
关于session.get的问题说明
session.get默认只会查询数据库以及session.dirty/session.deleted中的实例,不会包含session.new里未执行flush的实例。如果要让它查到刚添加的实例,需要先执行session.flush(),但这会把pending数据同步到数据库(未提交),会增加数据库交互次数,失去批量提交的优势,因此不推荐这种方式。
内容的提问来源于stack exchange,提问作者Gustavo Costa
相关产品推荐
相关产品推荐

