从SQL Server迁移海量数据到Parquet文件的优化方案咨询
关于SQL Server流式导出Parquet的方案及优化建议
一、无需加载全量DataTable的流式处理方案
现有方案的核心瓶颈是一次性把全表数据加载到DataTable,大表场景下必然内存溢出,直接用SqlDataReader做流式逐批读取即可,完全不需要经过DataTable中转:
- 把
GetData方法替换为返回SqlDataReader的实现,配置CommandBehavior.SequentialAccess参数让ADO.NET不缓存全量查询结果,改为逐行流式读取 - 每读取满一个行组的行数就直接写入Parquet,内存中仅保留当前行组的数据,不会留存全表内容
- 如需拆分多文件,可以按固定行数(例如每100万行一个文件)、或者SQL Server分页查询拆分任务,用多进程/多线程并行处理不同分页区间,大幅提升导出速度
核心逻辑修改参考:
static void StreamWriteParquet(string connectionString, string query, string outputDir, int rowsPerFile = 1000000, int rowGroupSize = 100000) { using var conn = new SqlConnection(connectionString); conn.Open(); using var cmd = new SqlCommand(query, conn); // 关键配置:流式读取,不缓存全量结果 using var reader = cmd.ExecuteReader(CommandBehavior.SequentialAccess); // 先读取数据源Schema生成Parquet字段定义 var fields = GenerateSchemaFromReader(reader); int currentFileRowCount = 0; int fileIndex = 0; ParquetWriter writer = null; Stream fileStream = null; IList[] columnBuffers = new IList[reader.FieldCount]; for (int i = 0; i < reader.FieldCount; i++) { var targetType = reader.GetFieldType(i); if (targetType == typeof(DateTime)) targetType = typeof(DateTimeOffset); var valueType = targetType.IsClass ? targetType : typeof(Nullable<>).MakeGenericType(targetType); columnBuffers[i] = (IList)Activator.CreateInstance(typeof(List<>).MakeGenericType(valueType)); } try { while (reader.Read()) { // 写满单文件行数限制则关闭当前文件,创建新文件 if (currentFileRowCount >= rowsPerFile) { WriteRowGroup(writer, columnBuffers, fields); writer.Dispose(); fileStream.Dispose(); currentFileRowCount = 0; fileIndex++; } // 新文件初始化 if (writer == null || currentFileRowCount == 0) { var outputPath = Path.Combine(outputDir, $"part_{fileIndex}.parquet"); fileStream = File.Open(outputPath, FileMode.Create, FileAccess.Write); writer = new ParquetWriter(new Schema(fields), fileStream); } // 读取单行列数据存入缓冲区 for (int i = 0; i < reader.FieldCount; i++) { if (reader.IsDBNull(i)) { columnBuffers[i].Add(null); } else { if (reader.GetFieldType(i) == typeof(DateTime)) { columnBuffers[i].Add(new DateTimeOffset(reader.GetDateTime(i))); } else { columnBuffers[i].Add(reader.GetValue(i)); } } } currentFileRowCount++; // 缓冲区满一个行组则写入Parquet if (currentFileRowCount % rowGroupSize == 0) { WriteRowGroup(writer, columnBuffers, fields); // 清空缓冲区等待下一批数据 foreach (var buf in columnBuffers) buf.Clear(); } } // 写入最后不足一个行组的剩余数据 if (columnBuffers[0].Count > 0) { WriteRowGroup(writer, columnBuffers, fields); } } finally { writer?.Dispose(); fileStream?.Dispose(); reader?.Dispose(); conn?.Close(); } } // 行组写入辅助方法,逻辑与原有代码对齐 static void WriteRowGroup(ParquetWriter writer, IList[] columnBuffers, Field[] fields) { using var rgw = writer.CreateRowGroup(); for (int i = 0; i < columnBuffers.Length; i++) { var valueType = columnBuffers[i].GetType().GetGenericArguments()[0]; var arr = Array.CreateInstance(valueType, columnBuffers[i].Count); columnBuffers[i].CopyTo(arr, 0); rgw.WriteColumn(new DataColumn(fields[i], arr)); } }
二、Parquet压缩率影响因素及优化方法
Parquet压缩率完全不只是由行组数量决定,影响优先级从高到低如下:
- 压缩算法选择:默认Snappy是速度优先,压缩率偏低;切换为ZSTD、Gzip可以大幅提升压缩率,通常ZSTD是速度和压缩率平衡的最优选择
- 数据排序:相同列的相近值连续存储时,字典编码和压缩算法的效率会明显提升,例如按日期、地域等维度排序后导出,压缩率甚至可以翻数倍
- 行组大小:行组太小会导致每个行组的字典冗余,拉低压缩率;行组太大则会影响后续读取时的过滤效率,通常单文件行组大小设为10万100万行,单文件大小控制在128MB1GB之间是最优区间
- 数据类型选择:尽量使用精度匹配的小类型,比如不要用字符串存储数字、不要用DateTime存储仅日期的值,数据占用空间越小压缩率越高
三、超大规模表(2000亿行)优化建议
2000亿行的规模单线程导出完全无法满足需求,按以下方案调整:
- 任务拆分:将大表按主键范围、日期分区拆分为多个独立查询任务,用多进程/分布式调度并行处理,每个任务处理1000万~1亿行数据,互不干扰
- 性能优化:原有代码中大量使用反射创建列表、转换数组,性能损耗非常高,可以预生成各列的类型转换委托,或者直接使用Parquet.Net内置的强类型
IEnumerable<T>写入能力,性能比反射方案高2~3倍 - 内存控制:行组大小不要超过100万,单个导出任务的内存占用可以控制在数百MB以内,避免频繁GC卡顿
- IO优化:导出的Parquet文件不要写入系统盘,使用高速SSD或者直接写入对象存储,避免IO瓶颈
内容的提问来源于stack exchange,提问作者Tyron78
相关产品推荐
相关产品推荐

