如何使用SQLAlchemy engine批量执行多条PostgreSQL语句
批量执行SQL语句优化方案
原有代码存在两个严重影响执行效率的问题:
- 每处理单条数据就重新创建一次数据库连接,连接的创建和销毁会产生大量不必要的开销
- 所有SQL语句单条提交执行,和数据库的IO交互次数和数据量成正比,执行效率极低
优化思路
- 将数据库连接的创建逻辑移到循环外层,整批数据处理共用同一个连接
- 先按操作类型批量收集SQL语句和对应参数,统一调用批量执行接口,减少IO交互次数
- 增加显式事务控制,保证同批次数据操作的原子性,避免部分执行成功部分失败的数据不一致问题
优化后代码
with engine.connect() as connection: # 开启事务 trans = connection.begin() try: for batch in batches: # batch为字典对象构成的列表,batches为列表构成的外层列表 # 分别收集两类操作的参数列表 delete_queries_one_params = [] delete_queries_two_params = [] update_query_params = [] # 缓存同批次共用的SQL模板,避免重复生成 delete_template_one = None delete_template_two = None update_template = None for obj in batch: # 遍历列表内的每个字典对象 overwrite_failure = context['retry_failure'] == True and obj['status'] in ['success', 'critical failure'] if overwrite_failure: # load_deletes返回转义后的SQL模板,仅首次生成即可 if not delete_template_one: delete_template_one, delete_template_two = load_deletes(obj, table_name, columns, unique_cols, should_overwrite, identifier) delete_queries_one_params.append(obj) delete_queries_two_params.append(obj) else: # load_updates返回转义后的SQL模板,仅首次生成即可 if not update_template: update_template = load_updates(obj, table_name, columns, unique_cols, should_overwrite) update_query_params.append(obj) # 批量执行所有同类型操作 if delete_queries_one_params: connection.execute(delete_template_one, delete_queries_one_params) if delete_queries_two_params: connection.execute(delete_template_two, delete_queries_two_params) if update_query_params: connection.execute(update_template, update_query_params) # 整批处理完成后统一提交 trans.commit() except Exception as e: # 出现异常统一回滚,保证数据一致性 trans.rollback() raise e
额外注意事项
如果使用的是SQLAlchemy 2.0及以上版本,批量执行参数列表需要用bindparams包装,或者直接调用execute方法传入参数列表即可,框架会自动处理批量执行逻辑。如果是其他数据库驱动,也可以对应替换为对应驱动的批量执行接口,整体逻辑不变。
内容的提问来源于stack exchange,提问作者HermSang
相关产品推荐
相关产品推荐

