diff --git a/parquet-column/src/main/java/org/apache/parquet/CorruptStatistics.java b/parquet-column/src/main/java/org/apache/parquet/CorruptStatistics.java index c5846f9efa..546b9bdf25 100644 --- a/parquet-column/src/main/java/org/apache/parquet/CorruptStatistics.java +++ b/parquet-column/src/main/java/org/apache/parquet/CorruptStatistics.java @@ -70,37 +70,67 @@ public static boolean shouldIgnoreStatistics(String createdBy, PrimitiveTypeName try { ParsedVersion version = VersionParser.parse(createdBy); + return shouldIgnoreStatistics(version, createdBy, columnType); + } catch (RuntimeException | VersionParseException e) { + // couldn't parse the created_by field, log what went wrong, don't trust the + // stats, but don't make this fatal. + warnParseErrorOnce(createdBy, e); + return true; + } + } - if (!"parquet-mr".equals(version.application)) { - // assume other applications don't have this bug - return false; - } + /** + * Decides if the statistics from a file should be ignored because they are potentially corrupt. + * Use this when the writer version has already been parsed to avoid redundant parsing. + * + * @param writerVersion the pre-parsed writer version, or {@code null} if unknown/unparseable + * @param createdBy the original created-by string from the file footer (used for logging) + * @param columnType the type of the column that this is checking + * @return true if the statistics may be invalid and should be ignored, false otherwise + */ + public static boolean shouldIgnoreStatistics( + ParsedVersion writerVersion, String createdBy, PrimitiveTypeName columnType) { - if (Strings.isNullOrEmpty(version.version)) { - warnOnce("Ignoring statistics because created_by did not contain a semver (see PARQUET-251): " - + createdBy); - return true; - } + if (columnType != PrimitiveTypeName.BINARY && columnType != PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY) { + return false; + } + + if (writerVersion == null) { + warnOnce("Ignoring statistics because created_by is null or empty! See PARQUET-251 and PARQUET-297"); + return true; + } - SemanticVersion semver = SemanticVersion.parse(version.version); + if (!"parquet-mr".equals(writerVersion.application)) { + return false; + } + + if (Strings.isNullOrEmpty(writerVersion.version)) { + warnOnce("Ignoring statistics because created_by did not contain a semver (see PARQUET-251): " + createdBy); + return true; + } - if (semver.compareTo(PARQUET_251_FIXED_VERSION) < 0 - && !(semver.compareTo(CDH_5_PARQUET_251_FIXED_START) >= 0 - && semver.compareTo(CDH_5_PARQUET_251_FIXED_END) < 0)) { - warnOnce("Ignoring statistics because this file was created prior to " - + PARQUET_251_FIXED_VERSION - + ", see PARQUET-251"); - return true; + if (!writerVersion.hasSemanticVersion()) { + try { + SemanticVersion.parse(writerVersion.version); + } catch (SemanticVersionParseException e) { + warnParseErrorOnce(createdBy, e); } + return true; + } - // this file was created after the fix - return false; - } catch (RuntimeException | SemanticVersionParseException | VersionParseException e) { - // couldn't parse the created_by field, log what went wrong, don't trust the stats, - // but don't make this fatal. - warnParseErrorOnce(createdBy, e); + SemanticVersion semver = writerVersion.getSemanticVersion(); + + if (semver.compareTo(PARQUET_251_FIXED_VERSION) < 0 + && !(semver.compareTo(CDH_5_PARQUET_251_FIXED_START) >= 0 + && semver.compareTo(CDH_5_PARQUET_251_FIXED_END) < 0)) { + warnOnce("Ignoring statistics because this file was created prior to " + + PARQUET_251_FIXED_VERSION + + ", see PARQUET-251"); return true; } + + // this file was created after the fix + return false; } private static void warnParseErrorOnce(String createdBy, Throwable e) { diff --git a/parquet-column/src/test/java/org/apache/parquet/CorruptStatisticsTest.java b/parquet-column/src/test/java/org/apache/parquet/CorruptStatisticsTest.java index eb8b0b4b44..ff950a252b 100644 --- a/parquet-column/src/test/java/org/apache/parquet/CorruptStatisticsTest.java +++ b/parquet-column/src/test/java/org/apache/parquet/CorruptStatisticsTest.java @@ -20,6 +20,7 @@ import static org.assertj.core.api.Assertions.assertThat; +import org.apache.parquet.VersionParser.ParsedVersion; import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName; import org.junit.jupiter.api.Test; @@ -129,6 +130,49 @@ public void testCorruptStatistics() { .isFalse(); } + @Test + public void testShouldIgnoreStatisticsWithParsedVersion() throws Exception { + String createdBy = "parquet-mr version 1.6.0 (build abc)"; + + assertThat(CorruptStatistics.shouldIgnoreStatistics(null, null, PrimitiveTypeName.BINARY)) + .isTrue(); + + assertThat(CorruptStatistics.shouldIgnoreStatistics(null, null, PrimitiveTypeName.INT32)) + .isFalse(); + + ParsedVersion impala = VersionParser.parse("impala version 1.2.0 (build abc)"); + assertThat(CorruptStatistics.shouldIgnoreStatistics( + impala, "impala version 1.2.0 (build abc)", PrimitiveTypeName.BINARY)) + .isFalse(); + + ParsedVersion corrupt = VersionParser.parse(createdBy); + assertThat(CorruptStatistics.shouldIgnoreStatistics(corrupt, createdBy, PrimitiveTypeName.BINARY)) + .isTrue(); + + ParsedVersion fixed = VersionParser.parse("parquet-mr version 1.8.0 (build abc)"); + assertThat(CorruptStatistics.shouldIgnoreStatistics( + fixed, "parquet-mr version 1.8.0 (build abc)", PrimitiveTypeName.BINARY)) + .isFalse(); + + ParsedVersion newer = VersionParser.parse("parquet-mr version 1.12.0 (build abc)"); + assertThat(CorruptStatistics.shouldIgnoreStatistics( + newer, "parquet-mr version 1.12.0 (build abc)", PrimitiveTypeName.BINARY)) + .isFalse(); + + // version field present but not a valid semantic version + ParsedVersion invalidSemver = new ParsedVersion("parquet-mr", "not-a-semver", "abc"); + assertThat(invalidSemver.hasSemanticVersion()).isFalse(); + assertThat(CorruptStatistics.shouldIgnoreStatistics( + invalidSemver, "parquet-mr version not-a-semver (build abc)", PrimitiveTypeName.BINARY)) + .isTrue(); + + // empty version field + ParsedVersion emptyVersion = new ParsedVersion("parquet-mr", "", "abc"); + assertThat(CorruptStatistics.shouldIgnoreStatistics( + emptyVersion, "parquet-mr version (build abc)", PrimitiveTypeName.BINARY)) + .isTrue(); + } + @Test public void testDistributionCorruptStatistics() { assertThat(CorruptStatistics.shouldIgnoreStatistics( diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java b/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java index 465516e48f..1f2fd6e4ad 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java @@ -47,6 +47,8 @@ import org.apache.parquet.CorruptStatistics; import org.apache.parquet.ParquetReadOptions; import org.apache.parquet.Preconditions; +import org.apache.parquet.VersionParser.ParsedVersion; +import org.apache.parquet.VersionParser.VersionParseException; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.EncodingStats; import org.apache.parquet.column.ParquetProperties; @@ -945,7 +947,6 @@ public static org.apache.parquet.column.statistics.Statistics fromParquetStatist // Visible for testing static org.apache.parquet.column.statistics.Statistics fromParquetStatisticsInternal( String createdBy, Statistics formatStats, PrimitiveType type, SortOrder typeSortOrder) { - // create stats object based on the column type org.apache.parquet.column.statistics.Statistics.Builder statsBuilder = org.apache.parquet.column.statistics.Statistics.getBuilderForReading(type); @@ -986,12 +987,65 @@ static org.apache.parquet.column.statistics.Statistics fromParquetStatisticsInte return statsBuilder.build(); } + // Visible for testing + static org.apache.parquet.column.statistics.Statistics fromParquetStatisticsInternal( + ParsedVersion writerVersion, + String createdBy, + Statistics formatStats, + PrimitiveType type, + SortOrder typeSortOrder) { + org.apache.parquet.column.statistics.Statistics.Builder statsBuilder = + org.apache.parquet.column.statistics.Statistics.getBuilderForReading(type); + + if (formatStats != null) { + // Use the new V2 min-max statistics over the former one if it is filled + if (formatStats.isSetMin_value() && formatStats.isSetMax_value()) { + byte[] min = formatStats.min_value.array(); + byte[] max = formatStats.max_value.array(); + if (isMinMaxStatsSupported(type) || Arrays.equals(min, max)) { + statsBuilder.withMin(min); + statsBuilder.withMax(max); + } + } else { + boolean isSet = formatStats.isSetMax() && formatStats.isSetMin(); + boolean maxEqualsMin = isSet ? Arrays.equals(formatStats.getMin(), formatStats.getMax()) : false; + boolean sortOrdersMatch = SortOrder.SIGNED == typeSortOrder; + // NOTE: See docs in CorruptStatistics for explanation of why this check is needed + // The sort order is checked to avoid returning min/max stats that are not + // valid with the type's sort order. In previous releases, all stats were + // aggregated using a signed byte-wise ordering, which isn't valid for all the + // types (e.g. strings, decimals etc.). + if (!CorruptStatistics.shouldIgnoreStatistics(writerVersion, createdBy, type.getPrimitiveTypeName()) + && (sortOrdersMatch || maxEqualsMin)) { + if (isSet) { + statsBuilder.withMin(formatStats.min.array()); + statsBuilder.withMax(formatStats.max.array()); + } + } + } + + if (formatStats.isSetNull_count()) { + statsBuilder.withNumNulls(formatStats.null_count); + } + if (formatStats.isSetNan_count()) { + statsBuilder.withNanCount(formatStats.getNan_count()); + } + } + return statsBuilder.build(); + } + public org.apache.parquet.column.statistics.Statistics fromParquetStatistics( String createdBy, Statistics statistics, PrimitiveType type) { SortOrder expectedOrder = overrideSortOrderToSigned(type) ? SortOrder.SIGNED : sortOrder(type); return fromParquetStatisticsInternal(createdBy, statistics, type, expectedOrder); } + public org.apache.parquet.column.statistics.Statistics fromParquetStatistics( + ParsedVersion writerVersion, String createdBy, Statistics statistics, PrimitiveType type) { + SortOrder expectedOrder = overrideSortOrderToSigned(type) ? SortOrder.SIGNED : sortOrder(type); + return fromParquetStatisticsInternal(writerVersion, createdBy, statistics, type, expectedOrder); + } + GeospatialStatistics toParquetGeospatialStatistics( org.apache.parquet.column.statistics.geospatial.GeospatialStatistics geospatialStatistics) { if (geospatialStatistics == null) { @@ -1837,6 +1891,28 @@ public ColumnChunkMetaData buildColumnChunkMetaData( fromParquetStatistics(metaData.geospatial_statistics, type)); } + public ColumnChunkMetaData buildColumnChunkMetaData( + ColumnMetaData metaData, + ColumnPath columnPath, + PrimitiveType type, + ParsedVersion writerVersion, + String createdBy) { + return ColumnChunkMetaData.get( + columnPath, + type, + fromFormatCodec(metaData.codec), + convertEncodingStats(metaData.getEncoding_stats()), + fromFormatEncodings(metaData.encodings), + fromParquetStatistics(writerVersion, createdBy, metaData.statistics, type), + metaData.data_page_offset, + metaData.dictionary_page_offset, + metaData.num_values, + metaData.total_compressed_size, + metaData.total_uncompressed_size, + fromParquetSizeStatistics(metaData.size_statistics, type), + fromParquetStatistics(metaData.geospatial_statistics, type)); + } + public ParquetMetadata fromParquetMetadata(FileMetaData parquetMetadata) throws IOException { return fromParquetMetadata(parquetMetadata, null, false); } @@ -1854,6 +1930,17 @@ public ParquetMetadata fromParquetMetadata( Map rowGroupToRowIndexOffsetMap) throws IOException { MessageType messageType = fromParquetSchema(parquetMetadata.getSchema(), parquetMetadata.getColumn_orders()); + org.apache.parquet.hadoop.metadata.FileMetaData fileMetaData = + buildFileMetaData(parquetMetadata, messageType, encryptedFooter, fileDecryptor); + String createdBy = fileMetaData.getCreatedBy(); + ParsedVersion writerVersion = null; + boolean useWriterVersion = false; + try { + writerVersion = fileMetaData.getWriterVersion(); + useWriterVersion = true; + } catch (VersionParseException e) { + // Fall back to String-based path which logs the parse error with full context + } List blocks = new ArrayList(); List row_groups = parquetMetadata.getRow_groups(); @@ -1930,13 +2017,13 @@ public ParquetMetadata fromParquetMetadata( } } - String createdBy = parquetMetadata.getCreated_by(); if (!lazyMetadataDecryption) { // full column metadata (with stats) is available - column = buildColumnChunkMetaData( - metaData, - columnPath, - messageType.getType(columnPath.toArray()).asPrimitiveType(), - createdBy); + PrimitiveType primitiveType = + messageType.getType(columnPath.toArray()).asPrimitiveType(); + column = useWriterVersion + ? buildColumnChunkMetaData( + metaData, columnPath, primitiveType, writerVersion, createdBy) + : buildColumnChunkMetaData(metaData, columnPath, primitiveType, createdBy); column.setRowGroupOrdinal(rowGroup.getOrdinal()); if (metaData.isSetBloom_filter_offset()) { column.setBloomFilterOffset(metaData.getBloom_filter_offset()); @@ -1975,6 +2062,15 @@ public ParquetMetadata fromParquetMetadata( blocks.add(blockMetaData); } } + return new ParquetMetadata(fileMetaData, blocks); + } + + private static org.apache.parquet.hadoop.metadata.FileMetaData buildFileMetaData( + FileMetaData parquetMetadata, + MessageType messageType, + boolean encryptedFooter, + InternalFileDecryptor fileDecryptor) { + String createdBy = parquetMetadata.getCreated_by(); Map keyValueMetaData = new HashMap(); List key_value_metadata = parquetMetadata.getKey_value_metadata(); if (key_value_metadata != null) { @@ -1990,10 +2086,8 @@ public ParquetMetadata fromParquetMetadata( } else { encryptionType = EncryptionType.UNENCRYPTED; } - return new ParquetMetadata( - new org.apache.parquet.hadoop.metadata.FileMetaData( - messageType, keyValueMetaData, parquetMetadata.getCreated_by(), encryptionType, fileDecryptor), - blocks); + return new org.apache.parquet.hadoop.metadata.FileMetaData( + messageType, keyValueMetaData, createdBy, encryptionType, fileDecryptor); } private static IndexReference toColumnIndexReference(ColumnChunk columnChunk) { diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/metadata/FileMetaData.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/metadata/FileMetaData.java index 4143dd805a..bb30bf39c7 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/metadata/FileMetaData.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/metadata/FileMetaData.java @@ -24,6 +24,10 @@ import java.io.Serializable; import java.util.Map; import java.util.Objects; +import org.apache.parquet.Strings; +import org.apache.parquet.VersionParser; +import org.apache.parquet.VersionParser.ParsedVersion; +import org.apache.parquet.VersionParser.VersionParseException; import org.apache.parquet.crypto.InternalFileDecryptor; import org.apache.parquet.schema.MessageType; @@ -39,9 +43,22 @@ public enum EncryptionType { ENCRYPTED_FOOTER } + private static final class WriterVersionResult { + private final ParsedVersion version; + private final VersionParseException versionParseException; + + static final WriterVersionResult MISSING = new WriterVersionResult(null, null); + + WriterVersionResult(ParsedVersion version, VersionParseException versionParseException) { + this.version = version; + this.versionParseException = versionParseException; + } + } + private final MessageType schema; private final Map keyValueMetaData; private final String createdBy; + private transient volatile WriterVersionResult writerVersionResult; private final InternalFileDecryptor fileDecryptor; private final EncryptionType encryptionType; @@ -118,4 +135,40 @@ public InternalFileDecryptor getFileDecryptor() { public EncryptionType getEncryptionType() { return encryptionType; } + + /** + * Returns the parsed writer version from the {@code createdBy} string. The result is + * computed lazily and cached. + * + * @return the parsed version, or {@code null} if {@code createdBy} is null or empty + * @throws VersionParseException if {@code createdBy} is present but cannot be parsed + */ + @JsonIgnore + public ParsedVersion getWriterVersion() throws VersionParseException { + WriterVersionResult result = writerVersionResult; + if (result == null) { + synchronized (this) { + result = writerVersionResult; + if (result == null) { + result = parseCreatedBy(); + writerVersionResult = result; + } + } + } + if (result.versionParseException != null) { + throw result.versionParseException; + } + return result.version; + } + + private WriterVersionResult parseCreatedBy() { + if (Strings.isNullOrEmpty(createdBy)) { + return WriterVersionResult.MISSING; + } + try { + return new WriterVersionResult(VersionParser.parse(createdBy), null); + } catch (VersionParseException e) { + return new WriterVersionResult(null, e); + } + } } diff --git a/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java b/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java index 4d361d6aa0..59053ab241 100644 --- a/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java +++ b/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java @@ -2279,4 +2279,35 @@ public void testColumnIndexNanCountsRoundTrip() { assertThat(roundTrip).isNotNull(); assertThat(roundTrip.getNanCounts()).containsExactly(1L, 0L, 0L); } + + @Test + public void testV2StatsDoNotTriggerCorruptStatisticsCheck() { + // Regression test: when V2 stats (min_value/max_value) are present, + // shouldIgnoreStatistics should NOT be evaluated. This ensures the + // one-shot warning is not consumed for columns that use V2 stats. + org.apache.parquet.format.Statistics formatStats = new org.apache.parquet.format.Statistics(); + formatStats.setMin_value(ByteBuffer.wrap(new byte[] {0})); + formatStats.setMax_value(ByteBuffer.wrap(new byte[] {1})); + formatStats.setNull_count(0); + + PrimitiveType binaryType = Types.required(PrimitiveTypeName.BINARY).named("test_binary"); + + // Use a corrupt writer version (pre-1.8.0) — if shouldIgnoreStatistics were eagerly + // evaluated, it would log a warning and consume the one-shot flag + org.apache.parquet.VersionParser.ParsedVersion corruptVersion = + new org.apache.parquet.VersionParser.ParsedVersion("parquet-mr", "1.6.0", "abc"); + + org.apache.parquet.column.statistics.Statistics result = + ParquetMetadataConverter.fromParquetStatisticsInternal( + corruptVersion, + "parquet-mr version 1.6.0 (build abc)", + formatStats, + binaryType, + ParquetMetadataConverter.SortOrder.SIGNED); + + // V2 stats should be used regardless of corrupt version — min/max should be set + assertThat(result.hasNonNullValue()).isTrue(); + assertThat(result.getMinBytes()).isEqualTo(new byte[] {0}); + assertThat(result.getMaxBytes()).isEqualTo(new byte[] {1}); + } } diff --git a/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/metadata/FileMetaDataTest.java b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/metadata/FileMetaDataTest.java new file mode 100644 index 0000000000..1a10a90606 --- /dev/null +++ b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/metadata/FileMetaDataTest.java @@ -0,0 +1,85 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.parquet.hadoop.metadata; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.util.Collections; +import org.apache.parquet.VersionParser.VersionParseException; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Type; +import org.junit.jupiter.api.Test; + +class FileMetaDataTest { + + private static final MessageType SCHEMA = new MessageType( + "test", new PrimitiveType(Type.Repetition.REQUIRED, PrimitiveType.PrimitiveTypeName.INT32, "id")); + + @Test + void validCreatedByIsParsed() throws Exception { + FileMetaData meta = + new FileMetaData(SCHEMA, Collections.emptyMap(), "parquet-mr version 1.12.0 (build abc123)"); + + assertThat(meta.getWriterVersion()).isNotNull(); + assertThat(meta.getWriterVersion().application).isEqualTo("parquet-mr"); + assertThat(meta.getWriterVersion().version).isEqualTo("1.12.0"); + assertThat(meta.getWriterVersion().appBuildHash).isEqualTo("abc123"); + } + + @Test + void nullCreatedByReturnsNullWriterVersion() throws Exception { + FileMetaData meta = new FileMetaData(SCHEMA, Collections.emptyMap(), null); + + assertThat(meta.getWriterVersion()).isNull(); + assertThat(meta.getCreatedBy()).isNull(); + } + + @Test + void emptyCreatedByReturnsNullWriterVersion() throws Exception { + FileMetaData meta = new FileMetaData(SCHEMA, Collections.emptyMap(), ""); + + assertThat(meta.getWriterVersion()).isNull(); + } + + @Test + void unparseableCreatedByThrowsVersionParseException() { + FileMetaData meta = new FileMetaData(SCHEMA, Collections.emptyMap(), "no-version-here"); + + assertThatThrownBy(meta::getWriterVersion).isInstanceOf(VersionParseException.class); + } + + @Test + void versionWithoutBuildHash() throws Exception { + FileMetaData meta = new FileMetaData(SCHEMA, Collections.emptyMap(), "parquet-mr version 1.8.0"); + + assertThat(meta.getWriterVersion()).isNotNull(); + assertThat(meta.getWriterVersion().application).isEqualTo("parquet-mr"); + assertThat(meta.getWriterVersion().version).isEqualTo("1.8.0"); + assertThat(meta.getWriterVersion().appBuildHash).isNull(); + } + + @Test + void writerVersionIsCached() throws Exception { + FileMetaData meta = new FileMetaData(SCHEMA, Collections.emptyMap(), "parquet-mr version 1.12.0 (build abc)"); + + assertThat(meta.getWriterVersion()).isSameAs(meta.getWriterVersion()); + } +}