如何加速从无索引SQL Server读取数据至DataFrame?
我正在使用Python和SQL将SQL Server源表的数据读取至DataFrame,再写入Snowflake目标表。目前采用OFFSET和FETCH NEXT分批读取,但读取速度极慢。这些表大多无聚集或非聚集索引,且当前无法修改底层表。
现有代码示例:
import polars as pl from snowflake.connector.pandas_tools import write_pandas # function to read data in from SQL Server in Batches def extract_batch(source_table, offset, limit): session = mssql_conn_session() try: sql = f""" SELECT * FROM {source_table} ORDER BY practicecode, monthyear, nsaccountid, providertype, vpgo OFFSET {offset} ROWS FETCH NEXT {limit} ROWS ONLY """ # reading into polars dataframe which is faster than pandas.read_sql df = pl.read_database(sql, session) except Exception as E: raise Exception(f"Failed while fetching data from {source_table} due to {E}") finally: session.close() return df # function that takes above df and writes to snowflake def write_to_snowflake(df, target_table): connection = session.connection().connection try: # write_pandas() will chunk + parallelize under the hood success, nchunks, nrows, _ = write_pandas( conn=connection, df=df.to_pandas(), table_name=target_table, chunk_size=1_000_000, parallel=8, quote_identifiers=False, use_logical_type=True ) if not success: raise RuntimeError(f"Failed while inserting data to {target_table}") connection.commit() utils.logger.info( f"Successfully inserted {nrows} rows to {target_table} in {nchunks} chunks." ) except Exception as E: connection.rollback() raise Exception(E) finally: connection.close()
我根据表大小和所需批次数循环执行代码,写入Snowflake的速度很快(write_pandas底层基于COPY INTO命令),但SQL Server读取速度对比之下极慢,希望得到加速读取的可行方案。
1. 替换OFFSET/FETCH NEXT为键集分页
OFFSET在无索引的表上会逐行扫描到指定偏移量,数据量越大效率越低。改用键集分页,基于当前ORDER BY的列作为分页键,每次记录最后一行的键值,下一批从该键值开始读取:
-- 假设上一批最后一行的键值为 practicecode='XXX', monthyear='YYYYMM', nsaccountid=123, providertype='A', vpgo='B' SELECT * FROM {source_table} WHERE (practicecode > 'XXX') OR (practicecode = 'XXX' AND monthyear > 'YYYYMM') OR (practicecode = 'XXX' AND monthyear = 'YYYYMM' AND nsaccountid > 123) OR (practicecode = 'XXX' AND monthyear = 'YYYYMM' AND nsaccountid = 123 AND providertype > 'A') OR (practicecode = 'XXX' AND monthyear = 'YYYYMM' AND nsaccountid = 123 AND providertype = 'A' AND vpgo > 'B') ORDER BY practicecode, monthyear, nsaccountid, providertype, vpgo FETCH NEXT {limit} ROWS ONLY
这种方式让SQL Server直接定位到起始位置,避免全表扫描。需要在每次读取后保存当前批次最后一行的键值,作为下一批的查询条件。
2. 增大单次读取的批次大小
如果当前批次太小,频繁建立连接和执行查询会累积额外开销。尝试将limit提升至50万-100万行(根据内存情况调整),减少循环次数,降低连接和查询的重复开销。
3. 并行读取SQL Server数据
利用Python多线程/多进程发起多个范围扫描查询,每个线程处理一个数据分片:
- 预先分析
practicecode等排序列的分布,拆分出多个互不重叠的区间 - 每个线程负责读取一个区间的数据
- 用Polars的
pl.concat高效合并所有线程返回的DataFrame
注意:控制线程数在2-4个起步,避免给SQL Server造成过大压力。
4. 使用SQL Server批量导出工具配合Polars读取
先用bcp或sqlcmd将数据导出为CSV/Parquet文件,再用Polars读取本地文件,速度远快于逐批查询数据库:
# bcp导出CSV示例 bcp "SELECT * FROM {source_table} ORDER BY practicecode,monthyear,nsaccountid,providertype,vpgo" queryout data.csv -S <服务器地址> -d <数据库名> -U <用户名> -P <密码> -c -t, -r\n
导出后可直接用pl.read_csv("data.csv")读取,甚至可以跳过DataFrame转换,直接用Snowflake的COPY INTO命令上传文件。
5. 优化Polars读取数据库的参数
- 安装
adbc-driver-mssql后,在pl.read_database中指定engine="adbc",ADBC驱动通常比传统DBAPI驱动性能更优 - 避免
SELECT *,只读取需要的列,减少数据传输量和内存占用
6. 复用数据库连接
当前代码每次读取批次都新建并关闭连接,累积开销很大。改为创建单个连接并复用,直到所有批次读取完成再关闭:
def extract_all_batches(source_table, limit): session = mssql_conn_session() try: last_keys = None while True: if last_keys is None: sql = f""" SELECT * FROM {source_table} ORDER BY practicecode,monthyear,nsaccountid,providertype,vpgo FETCH NEXT {limit} ROWS ONLY """ else: pc, my, ns, pt, vp = last_keys sql = f""" SELECT * FROM {source_table} WHERE (practicecode > '{pc}') OR (practicecode = '{pc}' AND monthyear > '{my}') OR (practicecode = '{pc}' AND monthyear = '{my}' AND nsaccountid > {ns}) OR (practicecode = '{pc}' AND monthyear = '{my}' AND nsaccountid = {ns} AND providertype > '{pt}') OR (practicecode = '{pc}' AND monthyear = '{my}' AND nsaccountid = {ns} AND providertype = '{pt}' AND vpgo > '{vp}') ORDER BY practicecode,monthyear,nsaccountid,providertype,vpgo FETCH NEXT {limit} ROWS ONLY """ df = pl.read_database(sql, session) if df.is_empty(): break yield df # 获取最后一行的键值 last_row = df.tail(1) last_keys = ( last_row.get_column("practicecode")[0], last_row.get_column("monthyear")[0], last_row.get_column("nsaccountid")[0], last_row.get_column("providertype")[0], last_row.get_column("vpgo")[0] ) finally: session.close()
这样仅建立一次连接,大幅减少连接建立的开销。
内容的提问来源于stack exchange,提问作者amnesic

