Skip to content
Draft
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 @@ -39,6 +39,7 @@
import org.apache.parquet.column.values.factory.ValuesWriterFactory;
import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridEncoder;
import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridValuesWriter;
import org.apache.parquet.hadoop.metadata.ColumnPath;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.parquet.schema.MessageType;

Expand Down Expand Up @@ -743,6 +744,18 @@ public Builder withBloomFilterEnabled(String columnPath, boolean enabled) {
return this;
}

/**
* Enable or disable the bloom filter for the specified component-based column path.
*
* @param columnPath the path components of the column
* @param enabled whether bloom filter shall be enabled
* @return this builder for method chaining
*/
public Builder withBloomFilterEnabled(String[] columnPath, boolean enabled) {
this.bloomFilterEnabled.withValue(ColumnPath.get(columnPath), enabled);
return this;
}

public Builder withRowGroupRowCountLimit(int rowCount) {
Preconditions.checkArgument(rowCount > 0, "Invalid row count limit for row groups: %s", rowCount);
rowGroupRowCountLimit = rowCount;
Expand Down Expand Up @@ -777,6 +790,18 @@ public Builder withStatisticsEnabled(String columnPath, boolean enabled) {
return this;
}

/**
* Enable or disable statistics for the specified component-based column path.
*
* @param columnPath the path components of the column
* @param enabled whether statistics shall be enabled
* @return this builder for method chaining
*/
public Builder withStatisticsEnabled(String[] columnPath, boolean enabled) {
this.statistics.withValue(ColumnPath.get(columnPath), enabled);
return this;
}

public Builder withStatisticsEnabled(boolean enabled) {
this.statistics.withDefaultValue(enabled);
this.statisticsEnabled = enabled;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,13 @@
import static org.apache.parquet.hadoop.metadata.CompressionCodecName.GZIP;
import static org.apache.parquet.hadoop.metadata.CompressionCodecName.SNAPPY;
import static org.apache.parquet.hadoop.metadata.CompressionCodecName.ZSTD;
import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT32;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

import org.apache.parquet.schema.MessageType;
import org.apache.parquet.schema.MessageTypeParser;
import org.apache.parquet.schema.Types;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;

Expand Down Expand Up @@ -81,6 +83,32 @@ public void columnCodec_otherColumnsUnaffected() {
assertThat(props.getColumnCodec(colC)).isNull();
}

@Test
public void componentPathsKeepDottedAndNestedColumnPropertiesDistinct() {
MessageType schema = Types.buildMessage()
.required(INT32)
.named("a.b")
.requiredGroup()
.required(INT32)
.named("b")
.named("a")
.named("msg");
ColumnDescriptor topLevel = schema.getColumnDescription(new String[] {"a.b"});
ColumnDescriptor nested = schema.getColumnDescription(new String[] {"a", "b"});

ParquetProperties props = ParquetProperties.builder()
.withBloomFilterEnabled(false)
.withBloomFilterEnabled(new String[] {"a.b"}, true)
.withStatisticsEnabled(false)
.withStatisticsEnabled(new String[] {"a.b"}, true)
.build();

assertThat(props.isBloomFilterEnabled(topLevel)).isTrue();
assertThat(props.isBloomFilterEnabled(nested)).isFalse();
assertThat(props.getStatisticsEnabled(topLevel)).isTrue();
assertThat(props.getStatisticsEnabled(nested)).isFalse();
}

@Test
public void columnLevel_setForColumn_returnsConfiguredLevel() {
ParquetProperties props =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,18 +19,29 @@

package org.apache.parquet.hadoop;

import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Base64;
import java.util.List;
import java.util.Map;
import java.util.function.BiConsumer;
import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import org.apache.hadoop.conf.Configuration;

/**
* Parses the specified key-values in the format of root.key#column.path from a {@link Configuration} object.
* Parses per-column values from a {@link Configuration}.
*
* <p>Legacy keys use {@code root.key#column.path}. Structured keys use {@code
* root.key.column-path#encoded.path}, where each UTF-8 path component is encoded independently so
* literal dots in field names remain distinct from path separators.
*/
class ColumnConfigParser {

private static final String COLUMN_PATH_SUFFIX = ".column-path";
private static final String EMPTY_PATH_COMPONENT = "~";

private static class ConfigHelper<T> {
private final String prefix;
private final Function<String, T> function;
Expand All @@ -51,14 +62,71 @@ public void processKey(String key) {
}
}

private static class ColumnPathConfigHelper<T> {
private final String prefix;
private final Function<String, T> function;
private final BiConsumer<String[], T> consumer;

private ColumnPathConfigHelper(String prefix, Function<String, T> function, BiConsumer<String[], T> consumer) {
this.prefix = prefix;
this.function = function;
this.consumer = consumer;
}

private void processKey(String key) {
if (key.startsWith(prefix)) {
String columnPath = key.substring(prefix.length());
T value = function.apply(key);
consumer.accept(decodeColumnPath(columnPath), value);
}
}
}

private final List<ConfigHelper<?>> helpers = new ArrayList<>();
private final List<ColumnPathConfigHelper<?>> columnPathHelpers = new ArrayList<>();

public <T> ColumnConfigParser withColumnConfig(
String rootKey, Function<String, T> function, BiConsumer<String, T> consumer) {
helpers.add(new ConfigHelper<T>(rootKey + '#', function, consumer));
return this;
}

public <T> ColumnConfigParser withColumnPathConfig(
String rootKey, Function<String, T> function, BiConsumer<String[], T> consumer) {
columnPathHelpers.add(new ColumnPathConfigHelper<T>(rootKey + COLUMN_PATH_SUFFIX + '#', function, consumer));
return this;
}

static String columnPathKey(String rootKey, String[] columnPath) {
return rootKey + COLUMN_PATH_SUFFIX + '#' + encodeColumnPath(columnPath);
}

private static String encodeColumnPath(String[] columnPath) {
return Stream.of(columnPath)
.map(ColumnConfigParser::encodePathComponent)
.collect(Collectors.joining("."));
}

private static String encodePathComponent(String component) {
if (component.isEmpty()) {
return EMPTY_PATH_COMPONENT;
}
return Base64.getUrlEncoder().withoutPadding().encodeToString(component.getBytes(StandardCharsets.UTF_8));
}

private static String[] decodeColumnPath(String columnPath) {
return Stream.of(columnPath.split("\\.", -1))
.map(ColumnConfigParser::decodePathComponent)
.toArray(String[]::new);
}

private static String decodePathComponent(String component) {
if (EMPTY_PATH_COMPONENT.equals(component)) {
return "";
}
return new String(Base64.getUrlDecoder().decode(component), StandardCharsets.UTF_8);
}

public void parseConfig(Configuration conf) {
for (Map.Entry<String, String> entry : conf) {
for (ConfigHelper<?> helper : helpers) {
Expand All @@ -68,5 +136,12 @@ public void parseConfig(Configuration conf) {
helper.processKey(entry.getKey());
}
}
// Structured-path overrides are applied last so an unambiguous setting wins over a legacy
// dot-string setting for the same property.
for (Map.Entry<String, String> entry : conf) {
for (ColumnPathConfigHelper<?> helper : columnPathHelpers) {
helper.processKey(entry.getKey());
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,14 @@ public static boolean getBloomFilterEnabled(Configuration conf) {
return conf.getBoolean(BLOOM_FILTER_ENABLED, DEFAULT_BLOOM_FILTER_ENABLED);
}

/**
* Enables or disables Bloom filters for an unambiguous component-based column path.
*/
public static void setBloomFilterEnabled(JobContext jobContext, String[] columnPath, boolean enabled) {
getConfiguration(jobContext)
.setBoolean(ColumnConfigParser.columnPathKey(BLOOM_FILTER_ENABLED, columnPath), enabled);
}

public static boolean getAdaptiveBloomFilterEnabled(Configuration conf) {
return conf.getBoolean(ADAPTIVE_BLOOM_FILTER_ENABLED, DEFAULT_ADAPTIVE_BLOOM_FILTER_ENABLED);
}
Expand Down Expand Up @@ -433,6 +441,14 @@ public static void setStatisticsEnabled(JobContext jobContext, String columnPath
getConfiguration(jobContext).set(STATISTICS_ENABLED + "#" + columnPath, String.valueOf(enabled));
}

/**
* Enables or disables statistics for an unambiguous component-based column path.
*/
public static void setStatisticsEnabled(JobContext jobContext, String[] columnPath, boolean enabled) {
getConfiguration(jobContext)
.setBoolean(ColumnConfigParser.columnPathKey(STATISTICS_ENABLED, columnPath), enabled);
}

public static boolean getStatisticsEnabled(Configuration conf, String columnPath) {
String columnSpecific = conf.get(STATISTICS_ENABLED + "#" + columnPath);
if (columnSpecific != null) {
Expand Down Expand Up @@ -555,6 +571,12 @@ public RecordWriter<Void, T> getRecordWriter(Configuration conf, Path file, Comp
STATISTICS_ENABLED,
key -> conf.getBoolean(key, ParquetProperties.DEFAULT_STATISTICS_ENABLED),
propsBuilder::withStatisticsEnabled)
.withColumnPathConfig(
BLOOM_FILTER_ENABLED, key -> conf.getBoolean(key, false), propsBuilder::withBloomFilterEnabled)
.withColumnPathConfig(
STATISTICS_ENABLED,
key -> conf.getBoolean(key, ParquetProperties.DEFAULT_STATISTICS_ENABLED),
propsBuilder::withStatisticsEnabled)
.withColumnConfig(
COMPRESSION,
key -> CompressionCodecName.fromConf(conf.get(key)),
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
/*
* 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;

import static org.assertj.core.api.Assertions.assertThat;

import java.util.HashMap;
import java.util.Map;
import org.apache.hadoop.conf.Configuration;
import org.junit.jupiter.api.Test;

public class TestColumnConfigParser {

@Test
public void structuredColumnPathsRoundTripWithoutCollisions() {
String rootKey = "parquet.test";
String[][] paths = {{"a.b"}, {"a", "b"}, {""}, {"münchen", "street.name"}};
Configuration conf = new Configuration(false);
for (int i = 0; i < paths.length; i++) {
conf.setInt(ColumnConfigParser.columnPathKey(rootKey, paths[i]), i);
}

Map<Integer, String[]> parsedPaths = new HashMap<>();
new ColumnConfigParser()
.withColumnPathConfig(
rootKey, key -> conf.getInt(key, -1), (path, value) -> parsedPaths.put(value, path))
.parseConfig(conf);

assertThat(parsedPaths).hasSize(paths.length);
assertThat(parsedPaths.get(0)).containsExactly("a.b");
assertThat(parsedPaths.get(1)).containsExactly("a", "b");
assertThat(parsedPaths.get(2)).containsExactly("");
assertThat(parsedPaths.get(3)).containsExactly("münchen", "street.name");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,56 @@ public void testParquetFileWithBloomFilter() throws IOException {
}
}

@Test
public void testStructuredStatisticsConfigKeepsCollidingDotStringPathsDistinct() throws Exception {
MessageType schema = Types.buildMessage()
.required(BINARY)
.as(stringType())
.named("a.b")
.requiredGroup()
.required(BINARY)
.as(stringType())
.named("b")
.named("a")
.named("msg");
Configuration conf = new Configuration();
GroupWriteSupport.setSchema(schema, conf);
conf.setBoolean(ParquetOutputFormat.STATISTICS_ENABLED, false);

// Component-based configuration can target ["a.b"] without also targeting ["a", "b"].
Job job = Job.getInstance(conf);
ParquetOutputFormat.setStatisticsEnabled(job, new String[] {"a.b"}, true);
conf = job.getConfiguration();

GroupFactory factory = new SimpleGroupFactory(schema);
Group group = factory.newGroup().append("a.b", "top-level");
group.addGroup("a").append("b", "nested");

Path path = newTempPath();
ParquetOutputFormat<Group> outputFormat = new ParquetOutputFormat<>(new GroupWriteSupport());
RecordWriter<Void, Group> writer = outputFormat.getRecordWriter(conf, path, UNCOMPRESSED);
try {
writer.write(null, group);
} finally {
writer.close(null);
}

try (ParquetFileReader reader = ParquetFileReader.open(HadoopInputFile.fromPath(path, conf))) {
BlockMetaData block = reader.getFooter().getBlocks().get(0);
ColumnChunkMetaData topLevelColumn = block.getColumns().stream()
.filter(column -> column.getPath().equals(ColumnPath.get("a.b")))
.findFirst()
.orElseThrow();
ColumnChunkMetaData nestedColumn = block.getColumns().stream()
.filter(column -> column.getPath().equals(ColumnPath.get("a", "b")))
.findFirst()
.orElseThrow();

assertThat(topLevelColumn.getStatistics().hasNonNullValue()).isTrue();
assertThat(nestedColumn.getStatistics().hasNonNullValue()).isFalse();
}
}

@Test
public void testParquetFileWithBloomFilterWithFpp() throws IOException {
int buildBloomFilterCount = 100000;
Expand Down