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

如何以流式方式将大型数组写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 00:02:04