如何用PySpark多连接将10亿条数据插入SQL Server与Oracle
实现方案:10亿条DataFrame多连接并行插入SQL Server与Oracle
核心思路
- 数据分块:将10亿条DataFrame拆分成若干小批次(建议每批次10万-100万条,根据数据库性能调整),避免内存溢出,同时提升单批次插入效率。
- 多连接并行:针对每个数据库创建多个独立连接,分配不同数据块并行插入;同时SQL Server和Oracle的插入任务可并行执行,最大化利用系统资源。
- 批量插入优化:结合数据库原生批量插入特性(如SQL Server的
fast_executemany、Oracle的数组绑定),配合Python驱动实现高效写入。
具体实现步骤(以Python为例)
1. 环境依赖
先安装必要的库:
pip install pandas pyodbc cx_Oracle sqlalchemy
2. 数据分块处理
拆分现有DataFrame为多个小批次:
import pandas as pd # 假设原始DataFrame为df,设置分块大小为50万条 chunk_size = 500000 chunks = [df[i:i+chunk_size] for i in range(0, len(df), chunk_size)]
3. 单数据库多连接并行插入
以SQL Server为例,用多线程实现并行插入:
import pyodbc from concurrent.futures import ThreadPoolExecutor # 每个线程创建独立的SQL Server连接 def get_sqlserver_conn(): conn_str = ( "DRIVER={ODBC Driver 17 for SQL Server};" "SERVER=你的SQLServer地址;" "DATABASE=目标库;" "UID=用户名;" "PWD=密码;" ) return pyodbc.connect(conn_str) # 单批次插入SQL Server def insert_sqlserver_chunk(chunk): conn = get_sqlserver_conn() cursor = conn.cursor() # 启用fast_executemany优化批量插入速度 cursor.fast_executemany = True insert_sql = """ INSERT INTO 目标表 (列1, 列2, 列3) VALUES (?, ?, ?) """ # 将DataFrame行转为元组列表 data = [tuple(row) for row in chunk.itertuples(index=False)] cursor.executemany(insert_sql, data) conn.commit() cursor.close() conn.close() # 启动多线程并行插入(max_workers根据CPU核心和数据库连接数限制调整) with ThreadPoolExecutor(max_workers=8) as executor: executor.map(insert_sqlserver_chunk, chunks)
4. 双数据库同时并行插入
结合多线程,同时启动SQL Server和Oracle的插入任务:
import cx_Oracle from concurrent.futures import ThreadPoolExecutor # 每个线程创建独立的Oracle连接 def get_oracle_conn(): dsn = cx_Oracle.makedsn("你的Oracle地址", 1521, service_name="你的服务名") return cx_Oracle.connect(user="用户名", password="密码", dsn=dsn) # 单批次插入Oracle def insert_oracle_chunk(chunk): conn = get_oracle_conn() cursor = conn.cursor() # 设置输入参数类型,启用数组绑定优化 cursor.setinputsizes(None, cx_Oracle.NCHAR, cx_Oracle.NUMBER) insert_sql = """ INSERT INTO 目标表 (列1, 列2, 列3) VALUES (:1, :2, :3) """ data = [tuple(row) for row in chunk.itertuples(index=False)] # 启用batcherrors允许跳过错误行继续插入 cursor.executemany(insert_sql, data, batcherrors=True) conn.commit() cursor.close() conn.close() # 同时启动双数据库的并行插入任务 def run_parallel_inserts(chunks): with ThreadPoolExecutor(max_workers=16) as executor: # 提交SQL Server所有分块任务 sql_tasks = [executor.submit(insert_sqlserver_chunk, chunk) for chunk in chunks] # 提交Oracle所有分块任务 oracle_tasks = [executor.submit(insert_oracle_chunk, chunk) for chunk in chunks] # 等待所有任务完成 for task in sql_tasks + oracle_tasks: task.result() # 执行并行插入 run_parallel_inserts(chunks)
关键优化点
- 连接池复用:用SQLAlchemy配置连接池,减少连接创建销毁的开销:
from sqlalchemy import create_engine sql_engine = create_engine("mssql+pyodbc://用户名:密码@地址/库?driver=ODBC+Driver+17+for+SQL+Server", pool_size=8, max_overflow=10) - 临时禁用索引/约束:插入前临时禁用目标表的非主键索引、外键约束,插入完成后再重建,大幅降低写入时的计算开销。
- 数据库配置调优:调整数据库的
max_connections参数,SQL Server可切换为简单恢复模式减少日志写入压力,Oracle可调整PGA_AGGREGATE_TARGET等内存参数。 - 分块大小调优:根据服务器内存、磁盘IO能力测试最优分块大小,过大易内存溢出,过小会增加连接开销。
内容的提问来源于stack exchange,提问作者Code Heaven
相关产品推荐
相关产品推荐

