Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@
package org.apache.parquet;

import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.parquet.SemanticVersion.SemanticVersionParseException;
import org.apache.parquet.VersionParser.ParsedVersion;
import org.apache.parquet.VersionParser.VersionParseException;
import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
Expand Down Expand Up @@ -70,37 +69,60 @@ public static boolean shouldIgnoreStatistics(String createdBy, PrimitiveTypeName

try {
ParsedVersion version = VersionParser.parse(createdBy);
return shouldIgnoreStatistics(version, columnType);
} catch (RuntimeException | VersionParseException e) {
warnParseErrorOnce(createdBy, e);
return true;
}
}

/**
* Decides if the statistics from a file should be ignored because they are potentially corrupt.
* Use this overload 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 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, PrimitiveTypeName columnType) {

if (!"parquet-mr".equals(version.application)) {
// assume other applications don't have this bug
return false;
}

if (Strings.isNullOrEmpty(version.version)) {
warnOnce("Ignoring statistics because created_by did not contain a semver (see PARQUET-251): "
+ createdBy);
return true;
}

SemanticVersion semver = SemanticVersion.parse(version.version);

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
if (columnType != PrimitiveTypeName.BINARY && columnType != PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY) {
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);
}

if (writerVersion == null) {
warnOnce("Ignoring statistics because created_by is null or empty! See PARQUET-251 and PARQUET-297");
return true;
}

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): "
+ writerVersion);
return true;
}

if (!writerVersion.hasSemanticVersion()) {
warnOnce("Ignoring statistics because created_by could not be parsed (see PARQUET-251): " + writerVersion);
return true;
}

SemanticVersion semver = writerVersion.getSemanticVersion();
Comment on lines +108 to +113

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ParsedVersion eagerly parses and caches the SemanticVersion in its constructor, so getSemanticVersion() avoids the redundant SemanticVersion.parse(version.version) that the String-based overload previously performed on every call. The left and right spikes in flame graph are for parsing SemanticVersion twice.

Image


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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -129,6 +130,42 @@ public void testCorruptStatistics() {
.isFalse();
}

@Test
public void testParsedVersionOverload() throws Exception {
assertThat(CorruptStatistics.shouldIgnoreStatistics((ParsedVersion) null, PrimitiveTypeName.BINARY))
.isTrue();

assertThat(CorruptStatistics.shouldIgnoreStatistics((ParsedVersion) null, PrimitiveTypeName.INT32))
.isFalse();

ParsedVersion impala = VersionParser.parse("impala version 1.2.0 (build abc)");
assertThat(CorruptStatistics.shouldIgnoreStatistics(impala, PrimitiveTypeName.BINARY))
.isFalse();

ParsedVersion corrupt = VersionParser.parse("parquet-mr version 1.6.0 (build abc)");
assertThat(CorruptStatistics.shouldIgnoreStatistics(corrupt, PrimitiveTypeName.BINARY))
.isTrue();

ParsedVersion fixed = VersionParser.parse("parquet-mr version 1.8.0 (build abc)");
assertThat(CorruptStatistics.shouldIgnoreStatistics(fixed, PrimitiveTypeName.BINARY))
.isFalse();

ParsedVersion newer = VersionParser.parse("parquet-mr version 1.12.0 (build abc)");
assertThat(CorruptStatistics.shouldIgnoreStatistics(newer, 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, PrimitiveTypeName.BINARY))
.isTrue();

// empty version field
ParsedVersion emptyVersion = new ParsedVersion("parquet-mr", "", "abc");
assertThat(CorruptStatistics.shouldIgnoreStatistics(emptyVersion, PrimitiveTypeName.BINARY))
.isTrue();
}

@Test
public void testDistributionCorruptStatistics() {
assertThat(CorruptStatistics.shouldIgnoreStatistics(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -945,7 +947,25 @@ 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
return fromParquetStatisticsInternal(
CorruptStatistics.shouldIgnoreStatistics(createdBy, type.getPrimitiveTypeName()),
formatStats,
type,
typeSortOrder);
}

// Visible for testing
static org.apache.parquet.column.statistics.Statistics fromParquetStatisticsInternal(
ParsedVersion writerVersion, Statistics formatStats, PrimitiveType type, SortOrder typeSortOrder) {
return fromParquetStatisticsInternal(
CorruptStatistics.shouldIgnoreStatistics(writerVersion, type.getPrimitiveTypeName()),
formatStats,
type,
typeSortOrder);
}

private static org.apache.parquet.column.statistics.Statistics fromParquetStatisticsInternal(
boolean shouldIgnoreStatistics, Statistics formatStats, PrimitiveType type, SortOrder typeSortOrder) {
org.apache.parquet.column.statistics.Statistics.Builder statsBuilder =
org.apache.parquet.column.statistics.Statistics.getBuilderForReading(type);

Expand All @@ -962,13 +982,7 @@ static org.apache.parquet.column.statistics.Statistics fromParquetStatisticsInte
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(createdBy, type.getPrimitiveTypeName())
&& (sortOrdersMatch || maxEqualsMin)) {
if (!shouldIgnoreStatistics && (sortOrdersMatch || maxEqualsMin)) {
if (isSet) {
statsBuilder.withMin(formatStats.min.array());
statsBuilder.withMax(formatStats.max.array());
Expand All @@ -992,6 +1006,12 @@ public org.apache.parquet.column.statistics.Statistics fromParquetStatistics(
return fromParquetStatisticsInternal(createdBy, statistics, type, expectedOrder);
}

public org.apache.parquet.column.statistics.Statistics fromParquetStatistics(
ParsedVersion writerVersion, Statistics statistics, PrimitiveType type) {
SortOrder expectedOrder = overrideSortOrderToSigned(type) ? SortOrder.SIGNED : sortOrder(type);
return fromParquetStatisticsInternal(writerVersion, statistics, type, expectedOrder);
}

GeospatialStatistics toParquetGeospatialStatistics(
org.apache.parquet.column.statistics.geospatial.GeospatialStatistics geospatialStatistics) {
if (geospatialStatistics == null) {
Expand Down Expand Up @@ -1837,6 +1857,24 @@ public ColumnChunkMetaData buildColumnChunkMetaData(
fromParquetStatistics(metaData.geospatial_statistics, type));
}

public ColumnChunkMetaData buildColumnChunkMetaData(
ColumnMetaData metaData, ColumnPath columnPath, PrimitiveType type, ParsedVersion writerVersion) {
return ColumnChunkMetaData.get(
columnPath,
type,
fromFormatCodec(metaData.codec),
convertEncodingStats(metaData.getEncoding_stats()),
fromFormatEncodings(metaData.encodings),
fromParquetStatistics(writerVersion, 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);
}
Expand All @@ -1854,6 +1892,17 @@ public ParquetMetadata fromParquetMetadata(
Map<RowGroup, Long> 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<BlockMetaData> blocks = new ArrayList<BlockMetaData>();
List<RowGroup> row_groups = parquetMetadata.getRow_groups();

Expand Down Expand Up @@ -1930,13 +1979,12 @@ 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)
: buildColumnChunkMetaData(metaData, columnPath, primitiveType, createdBy);
column.setRowGroupOrdinal(rowGroup.getOrdinal());
if (metaData.isSetBloom_filter_offset()) {
column.setBloomFilterOffset(metaData.getBloom_filter_offset());
Expand Down Expand Up @@ -1975,6 +2023,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<String, String> keyValueMetaData = new HashMap<String, String>();
List<KeyValue> key_value_metadata = parquetMetadata.getKey_value_metadata();
if (key_value_metadata != null) {
Expand All @@ -1990,10 +2047,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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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<String, String> keyValueMetaData;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This cached field is transient final. After Java deserialization it will stay null even when createdBy is valid. Recompute lazily in getWriterVersion(), or add serialization handling and tests.

private final String createdBy;
private transient volatile WriterVersionResult writerVersionResult;
private final InternalFileDecryptor fileDecryptor;
private final EncryptionType encryptionType;

Expand Down Expand Up @@ -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);
}
}
}
Loading