Parquet.Net 4.6.0异步写入Parquet文件后读取报错问题排查
问题原因与修正方案
你的Parquet文件写入代码核心问题是异步资源的生命周期管理错误,导致文件写入不完整(缺少Parquet必需的Footer元数据),因此读取时会报流结束相关错误。
具体错误点
- 错误将Task对象放入using块
using (var parquetWriterTask = ParquetWriter.CreateAsync(schema, fileStream))直接把CreateAsync返回的Task<ParquetWriter>放进using,using会立即Dispose这个Task对象,而非等待异步完成后得到的ParquetWriter实例,导致Writer的资源管理完全混乱。 - 使用
.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
相关产品推荐
相关产品推荐

