如何以流式方式将大型数组写入Parquet?
ParquetSharp流式写入IEnumerable<float?>避免内存不足
问题场景
我们使用ParquetSharp,定义的Schema如下:
new Column[] { new Column<string>("measurement_name"), new Column<DateTime?>("timestamp"), new Column<float?[]>("data") }
有时需要向data列写入约2.5亿条数据,目前必须将数据转换为float?[]才能写入,但这会引发内存不足错误。我们希望将MyDataset.Data属性改为IEnumerable<float?>类型,同时实现流式写入。虽然可以拆分数据集写入新RowGroup降低内存占用,但偶尔会遇到单个超大数据集的情况。
当前写入代码如下:
void WriteData(ParquetFileWriter writer1, IList<MyDataset> dataSets) { using var rg1 = writer1.AppendRowGroup(); using (var colWriter = rg1.NextColumn().LogicalWriter<string>()) { colWriter.WriteBatch(dataSets.Select(c => c.MeasurementName).ToArray()); } using (var colWriter = rg1.NextColumn().LogicalWriter<DateTime?>()) { colWriter.WriteBatch(dataSets.Select(c => c.DateTime).ToArray()); } using (var colWriter = rg1.NextColumn().LogicalWriter<float?[]>()) { colWriter.WriteBatch(dataSets.Select(c => c.Data).ToArray()); } }
解决方案
ParquetSharp中float?[]类型的列对应Parquet的LIST逻辑类型,底层分为数组长度和数组元素两个子列。我们可以直接操作这两个子列的Writer,实现流式写入IEnumerable<float?>数据,无需一次性将整个集合转换为数组。
1. 单RowGroup流式写入
适用于数据集行数不多,但部分行的data集合元素量极大的场景:
void WriteData(ParquetFileWriter writer1, IList<MyDataset> dataSets) { using var rg1 = writer1.AppendRowGroup(); // 写入measurement_name列 using (var colWriter = rg1.NextColumn().LogicalWriter<string>()) { colWriter.WriteBatch(dataSets.Select(c => c.MeasurementName).ToArray()); } // 写入timestamp列 using (var colWriter = rg1.NextColumn().LogicalWriter<DateTime?>()) { colWriter.WriteBatch(dataSets.Select(c => c.DateTime).ToArray()); } // 流式写入data列 var dataColumn = rg1.NextColumn(); // 获取数组长度的Writer(LIST类型的长度字段) using var lengthWriter = dataColumn.LogicalWriter<int>(); // 获取数组元素的Writer(LIST类型的元素字段) using var elementWriter = dataColumn.NextLogicalWriter<float?>(); var lengths = new List<int>(); var elementBatch = new List<float?>(); const int elementBatchSize = 100000; // 根据内存情况调整批次大小 foreach (var dataset in dataSets) { var data = dataset.Data; // 现在是IEnumerable<float?>类型 var elementCount = 0; foreach (var item in data) { elementBatch.Add(item); elementCount++; // 元素批次达到阈值时写入 if (elementBatch.Count >= elementBatchSize) { elementWriter.WriteBatch(elementBatch.ToArray()); elementBatch.Clear(); } } lengths.Add(elementCount); // 写入剩余的元素 if (elementBatch.Count > 0) { elementWriter.WriteBatch(elementBatch.ToArray()); elementBatch.Clear(); } } // 写入所有行的数组长度 lengthWriter.WriteBatch(lengths.ToArray()); }
2. 多RowGroup分批次写入
适用于总数据量极大(包括行数多或元素总量大)的场景,拆分数据集为多个小批次写入独立RowGroup,进一步控制内存占用:
void WriteData(ParquetFileWriter writer1, IEnumerable<MyDataset> dataSets) { const int rowGroupRowSize = 10000; // 每个RowGroup的行数,按需调整 var batch = new List<MyDataset>(rowGroupRowSize); foreach (var dataset in dataSets) { batch.Add(dataset); if (batch.Count >= rowGroupRowSize) { WriteSingleRowGroup(writer1, batch); batch.Clear(); } } // 写入剩余的行 if (batch.Count > 0) { WriteSingleRowGroup(writer1, batch); } } void WriteSingleRowGroup(ParquetFileWriter writer1, IList<MyDataset> batch) { using var rg = writer1.AppendRowGroup(); // 写入measurement_name列 using (var colWriter = rg.NextColumn().LogicalWriter<string>()) { colWriter.WriteBatch(batch.Select(c => c.MeasurementName).ToArray()); } // 写入timestamp列 using (var colWriter = rg.NextColumn().LogicalWriter<DateTime?>()) { colWriter.WriteBatch(batch.Select(c => c.DateTime).ToArray()); } // 流式写入data列 var dataColumn = rg.NextColumn(); using var lengthWriter = dataColumn.LogicalWriter<int>(); using var elementWriter = dataColumn.NextLogicalWriter<float?>(); var lengths = new List<int>(batch.Count); var elementBatch = new List<float?>(); const int elementBatchSize = 100000; foreach (var dataset in batch) { var elementCount = 0; foreach (var item in dataset.Data) { elementBatch.Add(item); elementCount++; if (elementBatch.Count >= elementBatchSize) { elementWriter.WriteBatch(elementBatch.ToArray()); elementBatch.Clear(); } } lengths.Add(elementCount); if (elementBatch.Count > 0) { elementWriter.WriteBatch(elementBatch.ToArray()); elementBatch.Clear(); } } lengthWriter.WriteBatch(lengths.ToArray()); }
注意事项
- 调整
elementBatchSize和rowGroupRowSize参数,平衡内存占用与写入性能。 - 确保
MyDataset.Data的IEnumerable<float?>支持多次枚举,若为单次枚举类型(如流式数据源),需保证枚举过程中数据不会丢失。 - 该方案依赖ParquetSharp对LIST类型的自动解析,需确保Schema定义
Column<float?[]>正确对应Parquet的LIST逻辑类型。
内容的提问来源于stack exchange,提问作者Niels Harremoes
相关产品推荐
相关产品推荐

