如何高效从SQL Server提取海量数据并保存为Parquet文件?
SQL Server海量数据转Parquet的性能优化方案
问题背景
现有SQL Server中数百万条数据,需转换为Parquet格式。当前使用pyodbc+pyarrow实现,按PERIOD字段(YYYYMM格式)分批查询,年度合并保存为example_YYYY.parquet文件,但每个文件处理耗时超20分钟,寻求优化方案。当前代码如下:
conn = pyodbc.connect('DRIVER={ODBC Driver 17 for SQL Server};' 'SERVER=SERVER;' 'DATABASE=DATABASE;' 'UID=USER;' 'PWD=PASS') cursor = conn.cursor() query = """SELECT DISTINCT [PERIOD] FROM [RIESGO_MODELO].[dbo].[MP_INCUMPLIMIENTO] ORDER BY [PERIOD] ASC ; """ cursor.execute(query) periods = [str(i[0]) for i in cursor.fetchall()] query_ = """SELECT [PERIOD] , [FIELD1] , ... , [FIELDN] FROM [RIESGO_MODELO].[dbo].[MP_INCUMPLIMIENTO] WHERE PERIOD = """ rows_cons = [] for period in periods: month = period[-2:] if month != '12': query = query_ + period + ';' rows = cursor.execute(query).fetchall() columns = [column[0] for column in cursor.description] rows_cons = rows_cons + rows else: query = query_ + period + ';' rows = cursor.execute(query).fetchall() columns = [column[0] for column in cursor.description] rows_cons = rows_cons + rows table = list(zip(*rows_cons)) pa_table = pa.table(table, names=columns) pq.write_table(pa_table, r'local_path\\example_{}.parquet'.format(period[:4])) rows_cons = [] del table del pa_table cursor.close() conn.close()
优化方案
1. 减少SQL查询次数:按年度批量查询
当前每个月执行一次SQL查询,年度需12次网络往返。改为直接按年度范围查询,将12次查询合并为1次,大幅降低网络IO和数据库执行开销:
# 提取所有不重复的年度 years = sorted(list(set([p[:4] for p in periods]))) for year in years: start_period = f"{year}01" end_period = f"{year}12" # 按年度范围执行参数化查询 query = """SELECT [PERIOD], [FIELD1], ..., [FIELDN] FROM [RIESGO_MODELO].[dbo].[MP_INCUMPLIMIENTO] WHERE PERIOD BETWEEN ? AND ? """ cursor.execute(query, (start_period, end_period)) # 后续处理年度数据
2. 替换fetchall()为分块读取,避免内存过载
fetchall()会一次性将所有数据加载到内存,百万级数据会造成内存压力且拖慢速度。改用fetchmany()分块读取,逐步构建Arrow表:
chunk_size = 100000 # 根据内存配置调整 pa_tables = [] columns = None while True: rows = cursor.fetchmany(chunk_size) if not rows: break if not columns: columns = [col[0] for col in cursor.description] # 将当前数据块转为Arrow表 chunk_table = pa.table(list(zip(*rows)), names=columns) pa_tables.append(chunk_table) # 合并所有数据块并写入Parquet combined_table = pa.concat_tables(pa_tables) pq.write_table(combined_table, f'local_path\\example_{year}.parquet')
3. 使用Pandas中转,简化流程并提升效率
Pandas对SQL读取和Parquet写入的优化更成熟,结合pyarrow引擎可大幅简化代码并提升性能:
import pandas as pd for year in years: start_period = f"{year}01" end_period = f"{year}12" query = """SELECT [PERIOD], [FIELD1], ..., [FIELDN] FROM [RIESGO_MODELO].[dbo].[MP_INCUMPLIMIENTO] WHERE PERIOD BETWEEN ? AND ? """ # 分块读取SQL数据,避免内存溢出 df_chunks = pd.read_sql(query, conn, params=(start_period, end_period), chunksize=100000) output_path = f'local_path\\example_{year}.parquet' # 分块写入Parquet for i, chunk in enumerate(df_chunks): if i == 0: chunk.to_parquet(output_path, engine='pyarrow', compression='snappy') else: chunk.to_parquet(output_path, engine='pyarrow', compression='snappy', mode='append')
4. 优化Parquet写入参数
调整Parquet写入配置,平衡压缩速度和文件大小:
- 使用
compression='snappy':比默认gzip压缩速度快数倍,压缩率满足大部分场景需求 - 匹配分块大小设置
row_group_size:提升读写效率 - 关闭不必要的统计信息:
write_statistics=False(无需统计查询时启用)
示例:
pq.write_table(combined_table, f'local_path\\example_{year}.parquet', compression='snappy', row_group_size=100000, write_statistics=False)
5. 数据库端优化
- 给
PERIOD字段添加非聚集索引:加速WHERE条件过滤,减少数据库全表扫描时间 - 检查查询执行计划:确认SQL是否走索引,避免低效扫描
- 只查询必要字段:移除SELECT语句中不需要的字段,减少数据传输量
6. 连接与SQL语句优化
- 使用参数化查询:避免字符串拼接SQL,防止注入同时提升数据库缓存效率
- 启用连接池:减少重复创建连接的开销:
conn = pyodbc.connect('DRIVER={ODBC Driver 17 for SQL Server};' 'SERVER=SERVER;' 'DATABASE=DATABASE;' 'UID=USER;' 'PWD=PASS;' 'MaxPoolSize=10;') # 设置连接池大小
7. 硬件与环境优化
- 使用SSD磁盘存储Parquet文件:大幅提升磁盘IO速度
- 分配足够内存:避免Python进程使用虚拟内存(swap)导致速度骤降
- 多进程并行处理:数据库支持并发时,用
multiprocessing同时处理多个年度数据,注意控制并发数避免数据库压力过大
内容的提问来源于stack exchange,提问作者imatiasmb
相关产品推荐
相关产品推荐

