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

Parquet.Net 4.6.0异步写入Parquet文件后读取报错问题排查

问题原因与修正方案

你的Parquet文件写入代码核心问题是异步资源的生命周期管理错误,导致文件写入不完整(缺少Parquet必需的Footer元数据),因此读取时会报流结束相关错误。

具体错误点

  1. 错误将Task对象放入using块
    using (var parquetWriterTask = ParquetWriter.CreateAsync(schema, fileStream)) 直接把CreateAsync返回的Task<ParquetWriter>放进using,using会立即Dispose这个Task对象,而非等待异步完成后得到的ParquetWriter实例,导致Writer的资源管理完全混乱。
  2. 使用.Wait()阻塞线程,异步操作未完全完成
    .Wait()会强制阻塞线程,但无法保证异步写入操作(包括Parquet文件的收尾工作)完全完成就进入后续的资源释放流程,最终文件缺少关键的Footer信息,变成不合法的Parquet文件。

修正后的异步写入代码

async Task WriteParquetFileAsync()
{
    string file = @"c:\temp\test.parquet";
    var dataFields = new DataField[2];
    dataFields[0] = new DataField("dtUTC", DataType.Int64);
    dataFields[1] = new DataField("val", DataType.Double);
    var schema = new ParquetSchema(dataFields);
    var dtUTC = new long[3];
    var val = new double[3];

    using (Stream fileStream = System.IO.File.OpenWrite(file))
    {
        // 先await获取实际的ParquetWriter实例,再放入using管理生命周期
        using (var parquetWriter = await ParquetWriter.CreateAsync(schema, fileStream))
        {
            using (ParquetRowGroupWriter groupWriter = parquetWriter.CreateRowGroup())
            {
                var col0 = Array.CreateInstance(typeof(long), dtUTC.Length);
                for(int i=0;i< dtUTC.Length;i++) col0.SetValue(dtUTC[i], i);
                // 使用await替代Wait(),确保异步写入完成
                await groupWriter.WriteColumnAsync(new Parquet.Data.DataColumn(dataFields[0], col0));
                
                var col1 = Array.CreateInstance(typeof(double), val.Length);
                for (int i = 0; i < val.Length; i++) col1.SetValue(val[i], i);
                await groupWriter.WriteColumnAsync(new Parquet.Data.DataColumn(dataFields[1], col1));
            }
            // ParquetWriter的Dispose会自动写入文件Footer,确保文件结构完整
        }
    }
}

补充:读取代码也应同步改为正确异步方式

async Task ReadParquetFileAsync(string fileName)
{
    using (Stream file = System.IO.File.OpenRead(fileName))
    {
        using (var reader = await ParquetReader.CreateAsync(file))
        {
            if (reader.RowGroupCount != 1) throw new Exception("reader.RowGroupCount = " + reader.RowGroupCount);
            var dataFields = reader.Schema.GetDataFields();
            using (var gr = reader.OpenRowGroupReader(0))
            {
                // 若getCols有异步版本,需用await调用
                Parquet.Data.DataColumn[] columns = await parquet.getColsAsync(gr, dataFields);
                Console.WriteLine($"\tGroup 0: columns.Length={columns.Length}");
            }
        }
    }
}

关键说明

Parquet文件的合法性依赖末尾的Footer元数据,ParquetWriter的Dispose方法会自动完成Footer的写入。只有通过await获取Writer实例并正确用using管理其生命周期,才能确保所有写入操作(包括Footer)完全完成,生成合法的Parquet文件。

内容的提问来源于stack exchange,提问作者mcmillab

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 03:24:53