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

批量提交时如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 07:47:41