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

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. 后续排查步骤

  1. 运行修正后的代码,查看打印的异常信息,定位具体失败原因(比如数据约束冲突、字段格式错误等)
  2. 查看PostgreSQL数据库日志,确认是否有数据库层面的报错或警告
  3. 先用小批量数据(比如100条)测试,验证修正后的代码能正常全部插入

内容的提问来源于stack exchange,提问作者Marvin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 05:07:29