ThreadPoolExecutor多线程插入PostgreSQL:无报错但数据插入不全
问题排查与修复方案
1. 先修正代码中的明显错误
你的代码存在几个直接导致任务静默失败的低级问题:
- 循环语法错误:
for future concurrent. futures.as_completed(futures)缺少in,正确写法是for future in concurrent.futures.as_completed(futures) - 变量名不匹配:
Session.add(one_politician)里的one_politician应该是one_new_entry - Session实例使用混乱:你已经获取了
session = Session(),后续操作应该用这个实例,而非直接调用Session.add()、Session.commit()这类类方法
修正后的核心代码:
engine = create_engine(DATABASE_URI) Session = scoped_session(sessionmaker(bind=engine)) def write_entry(information): session = Session() try: one_new_entry = ATypeOfEntry(a_field=information) session.add(one_new_entry) session.commit() except Exception as e: session.rollback() # 打印异常,明确失败原因 print(f"插入失败: {str(e)}") raise finally: Session.remove() if __name__ == "__main__": list_of_information = get_list_of_information() # 注意变量名统一,原代码中list_of_members_and_wikidata应为list_of_information with tqdm(total=len(list_of_information)) as pbar: with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor: futures = [executor.submit(write_entry, info) for info in list_of_information] for future in concurrent.futures.as_completed(futures): # 必须捕获future中的异常,否则失败任务会被隐藏 try: future.result() except Exception as e: print(f"任务执行异常: {str(e)}") pbar.update()
2. 核心问题分析
- 静默异常被忽略:ThreadPoolExecutor中任务抛出异常时,不会主动触发程序报错,只有调用
future.result()才会暴露。你之前的代码没有做这一步,导致任务失败但进度条仍正常更新,看起来全部完成。 - scoped_session使用错误:scoped_session是线程局部的,每个线程会绑定独立session,但你获取实例后又用类方法操作,导致session上下文混乱,引发未预期的提交失败。
3. 进一步优化建议
- 用批量插入替代单条插入:PostgreSQL对批量插入的支持远好于多线程单条插入,效率提升更明显。可以拆分数据批次处理:
def write_batch(batch_info): session = Session() try: entries = [ATypeOfEntry(a_field=info) for info in batch_info] session.bulk_save_objects(entries) session.commit() except Exception as e: session.rollback() print(f"批量插入失败: {str(e)}") raise finally: Session.remove() # 主函数中拆分批次 batch_size = 100 batches = [list_of_information[i:i+batch_size] for i in range(0, len(list_of_information), batch_size)] with tqdm(total=len(batches)) as pbar: with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor: futures = [executor.submit(write_batch, batch) for batch in batches] for future in concurrent.futures.as_completed(futures): try: future.result() except Exception as e: print(f"批次任务异常: {str(e)}") pbar.update()
- 调整连接池配置:创建engine时设置足够的连接池大小,避免线程因获取不到连接而阻塞:
engine = create_engine(DATABASE_URI, pool_size=10, max_overflow=20)
4. 后续排查步骤
- 运行修正后的代码,查看打印的异常信息,定位具体失败原因(比如数据约束冲突、字段格式错误等)
- 查看PostgreSQL数据库日志,确认是否有数据库层面的报错或警告
- 先用小批量数据(比如100条)测试,验证修正后的代码能正常全部插入
内容的提问来源于stack exchange,提问作者Marvin
相关产品推荐
相关产品推荐

