diff --git a/src/Parquet.Test/ParquetReaderOnTestFilesTest.cs b/src/Parquet.Test/ParquetReaderOnTestFilesTest.cs index c26256d6..f503a709 100644 --- a/src/Parquet.Test/ParquetReaderOnTestFilesTest.cs +++ b/src/Parquet.Test/ParquetReaderOnTestFilesTest.cs @@ -397,7 +397,6 @@ public async Task PyArrow23Async() { Assert.Equal(39_366, rs.Data.Count); } - [Fact] public async Task UnshreddedVariantAsync() { await using Stream s = OpenTestFile("variant_unshredded.parquet"); @@ -405,4 +404,18 @@ public async Task UnshreddedVariantAsync() { Assert.NotNull(r.Schema); } + [Fact] + public async Task AllNullColumnPyArrowV25() { + using Stream s = OpenTestFile("all_null_column_pyarrow_v25.parquet"); + await using ParquetReader r = await ParquetReader.CreateAsync(s); + using ParquetRowGroupReader groupReader = r.OpenRowGroupReader(0); + DataField[] fs = r.Schema.GetDataFields(); + double?[] data = await ReadNullableValuesAsync(groupReader, fs[0]); + Assert.Equal(46, data.Length); + Assert.All(data, d => Assert.Null(d)); + + data = await ReadNullableValuesAsync(groupReader, fs[1]); + Assert.Equal(46, data.Length); + Assert.All(data, d => Assert.Equal(0, d)); + } } diff --git a/src/Parquet.Test/data/all_null_column_pyarrow_v25.parquet b/src/Parquet.Test/data/all_null_column_pyarrow_v25.parquet new file mode 100644 index 00000000..6d4ddabb Binary files /dev/null and b/src/Parquet.Test/data/all_null_column_pyarrow_v25.parquet differ diff --git a/src/Parquet/File/DataColumnReader.cs b/src/Parquet/File/DataColumnReader.cs index b044c879..ebee9491 100644 --- a/src/Parquet/File/DataColumnReader.cs +++ b/src/Parquet/File/DataColumnReader.cs @@ -66,11 +66,11 @@ public async ValueTask ReadAsync(ReadingColumn rc, CancellationToken cance if(_stats?.NullCount != null) definedValuesCount -= (int)_stats.NullCount.Value; - //using var pc = new PackedColumn(_dataField, totalValuesInChunk, definedValuesCount); long fileOffset = GetFileOffset(); _inputStream.Seek(fileOffset, SeekOrigin.Begin); - while(rc.ValuesRead < totalValuesInChunk) { + bool allNullColumnProcessed = false; + while(rc.ValuesRead < totalValuesInChunk && !allNullColumnProcessed) { PageHeader ph = PageHeader.Read(new ThriftCompactProtocolReader(_inputStream)); switch(ph.Type) { @@ -79,6 +79,7 @@ public async ValueTask ReadAsync(ReadingColumn rc, CancellationToken cance break; case PageType.DATA_PAGE: await ReadDataPageV1Async(ph, rc, cancellationToken); + allNullColumnProcessed = definedValuesCount == 0; break; case PageType.DATA_PAGE_V2: await ReadDataPageV2Async(ph, rc, totalValuesInChunk, cancellationToken); @@ -122,12 +123,17 @@ private long GetFileOffset() => private async ValueTask ReadDataPageV1Async(PageHeader ph, ReadingColumn rc, CancellationToken cancellationToken) where T : struct { using IMemoryOwner bytes = await ReadPageDataAsync(ph); + int allValueCount = (int)_thriftColumnChunk.MetaData!.NumValues; if(ph.DataPageHeader == null) { - throw new ParquetException($"column '{_dataField.Path}' is missing data page header, file is corrupt"); + if (allValueCount != (int?)_stats?.NullCount) { + throw new ParquetException($"column '{_dataField.Path}' is missing data page header, file is corrupt"); + } + // all values are meant to be null; mark as read and return. + rc.MarkValuesRead(allValueCount); + return; } int dataUsed = 0; - int allValueCount = (int)_thriftColumnChunk.MetaData!.NumValues; int pageValueCount = ph.DataPageHeader.NumValues; if(_dataField.MaxRepetitionLevel > 0) {