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

如何高效从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 08:57:03