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

Python Pandas对接PostgreSQL:优化循环提升数据写入效率

PostgreSQL分层DataFrame写入优化方案

针对必须严格按客户→联系人→账户→产品层级顺序写入的场景,以下是提升10万+数据写入效率的优化方案:

1. 事务批量提交,降低数据库IO开销

默认psycopg2的单条execute隐式提交会产生大量网络IO,改为按客户批次提交事务,减少提交次数:

# 初始化时关闭自动提交
conn.autocommit = False
batch_size = 500  # 可根据服务器性能调整
counter = 0

try:
    for customer in customer_df.index.to_list():
        # 原有的客户层级所有写入逻辑...
        
        counter += 1
        if counter % batch_size == 0:
            conn.commit()
            print(f"已提交{counter}个客户的关联数据")
    # 提交剩余未批量的数据
    conn.commit()
except Exception as e:
    conn.rollback()
    raise e

2. 预分组DataFrame,减少索引扫描开销

原代码多次通过loc[[customer]]查询多层索引数据,重复触发索引扫描。提前按客户维度预分组,直接读取分组结果:

# 预执行分组,仅运行一次
contact_groups = contact_df.groupby(level=0)
account_groups = accounts_df.groupby(level=0)
product_groups = products_df.groupby(level=0)

# 循环中直接调用分组数据
for customer in customer_df.index.to_list():
    # 处理联系人数据
    if customer in contact_groups.groups:
        contact_group = contact_groups.get_group(customer)
        for idx, row in contact_group.iterrows():
            # 直接传入Series.values,无需转list
            cursor.execute(contact_sql, row.values)
            # 处理contact_df2同理
            cursor.execute(contact_sql_2, contact_df2.loc[idx].values)

同时简化generate_values_for_insert_statement函数,避免冗余类型判断:

def generate_values_for_insert_statement(df, index):
    row_data = df.loc[index]
    return [row_data.tolist()] if isinstance(row_data, pd.Series) else row_data.values.tolist()

3. 同层级数据批量写入:临时表+copy_from

同一客户下的同类型数据(如多个联系人、账户),用copy_from批量写入临时表,再一次性导入正式表,速度比逐行execute快10-100倍:

# 示例:批量写入同一客户的所有联系人数据
if customer in contact_groups.groups:
    contact_group = contact_groups.get_group(customer)
    # 创建临时表,结构与正式表一致,事务提交后自动销毁
    cursor.execute("CREATE TEMP TABLE temp_contact (LIKE contact_table INCLUDING ALL) ON COMMIT DROP;")
    # 将DataFrame转为Tab分隔的内存缓冲区
    from io import StringIO
    buffer = StringIO()
    contact_group.to_csv(buffer, sep='\t', header=False, index=False)
    buffer.seek(0)
    # 批量拷贝到临时表
    cursor.copy_from(buffer, 'temp_contact', sep='\t', columns=contact_group.columns.tolist())
    # 从临时表导入正式表
    cursor.execute("INSERT INTO contact_table SELECT * FROM temp_contact;")

4. 减少数据类型转换开销

psycopg2原生支持numpy数组和pandas Series作为参数,无需手动转为list,直接传入即可减少转换耗时:

# 原代码:
values = customer_df.loc[customer].to_list()
cursor.execute(customer_sql, values)
# 优化后:
cursor.execute(customer_sql, customer_df.loc[customer].values)

5. 连接池复用,避免重复建立连接

对于频繁运行的任务,使用连接池管理数据库连接,减少连接建立/销毁的开销:

from psycopg2 import pool

# 初始化连接池
connection_pool = pool.SimpleConnectionPool(
    minconn=1,
    maxconn=10,
    dbname="your_db",
    user="your_user",
    password="your_pass",
    host="your_host"
)

# 从连接池获取连接
conn = connection_pool.getconn()
cursor = conn.cursor()

# 执行所有写入操作...

# 归还连接到池
cursor.close()
connection_pool.putconn(conn)

# 程序结束时关闭连接池
connection_pool.closeall()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 04:17:03