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,17 +19,17 @@

package org.apache.parquet;

import static org.apache.parquet.conf.ParquetInputFilters.getFilter;
import static org.apache.parquet.conf.ParquetInputProperties.BLOOM_FILTERING_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.COLUMN_INDEX_FILTERING_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.DICTIONARY_FILTERING_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.HADOOP_VECTORED_IO_DEFAULT;
import static org.apache.parquet.conf.ParquetInputProperties.HADOOP_VECTORED_IO_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.OFF_HEAP_DECRYPT_BUFFER_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.PAGE_VERIFY_CHECKSUM_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.RECORD_FILTERING_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.STATS_FILTERING_ENABLED;
import static org.apache.parquet.format.converter.ParquetMetadataConverter.NO_FILTER;
import static org.apache.parquet.hadoop.ParquetInputFormat.BLOOM_FILTERING_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.COLUMN_INDEX_FILTERING_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.DICTIONARY_FILTERING_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.HADOOP_VECTORED_IO_DEFAULT;
import static org.apache.parquet.hadoop.ParquetInputFormat.HADOOP_VECTORED_IO_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.OFF_HEAP_DECRYPT_BUFFER_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.PAGE_VERIFY_CHECKSUM_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.RECORD_FILTERING_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.STATS_FILTERING_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.getFilter;
import static org.apache.parquet.hadoop.UnmaterializableRecordCounter.BAD_RECORD_THRESHOLD_CONF_KEY;

import java.util.Collections;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
/*
* 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.conf;

import java.io.IOException;
import org.apache.parquet.filter.UnboundRecordFilter;
import org.apache.parquet.filter2.compat.FilterCompat;
import org.apache.parquet.filter2.compat.FilterCompat.Filter;
import org.apache.parquet.filter2.predicate.FilterPredicate;
import org.apache.parquet.hadoop.BadConfigurationException;
import org.apache.parquet.hadoop.util.ConfigurationUtil;
import org.apache.parquet.hadoop.util.SerializationUtil;

/**
* Resolves the filter used to read Parquet files from a {@link ParquetConfiguration}.
* This class is decoupled from the legacy Hadoop {@code Configuration} and
* {@code InputFormat} machinery, so that the filter resolution logic can be used
* without loading {@code FileInputFormat} and its transitive dependencies.
*/
public final class ParquetInputFilters {

private ParquetInputFilters() {}

/**
* @param configuration a configuration
* @return an unbound record filter class
*/
public static Class<?> getUnboundRecordFilter(ParquetConfiguration configuration) {
return ConfigurationUtil.getClassFromConfig(
configuration, ParquetInputProperties.UNBOUND_RECORD_FILTER, UnboundRecordFilter.class);
}

/**
* @param configuration a configuration
* @return the filter predicate stored in the configuration, or null if none is set
*/
public static FilterPredicate getFilterPredicate(ParquetConfiguration configuration) {
try {
return SerializationUtil.readObjectFromConfAsBase64(ParquetInputProperties.FILTER_PREDICATE, configuration);
} catch (IOException e) {
throw new RuntimeException(e);
}
}

private static UnboundRecordFilter getUnboundRecordFilterInstance(ParquetConfiguration configuration) {
Class<?> clazz = getUnboundRecordFilter(configuration);
if (clazz == null) {
return null;
}
try {
return (UnboundRecordFilter) clazz.newInstance();
} catch (InstantiationException | IllegalAccessException e) {
throw new BadConfigurationException("could not instantiate unbound record filter class", e);
}
}

/**
* Returns a non-null Filter, which is a wrapper around either a
* FilterPredicate, an UnboundRecordFilter, or a no-op filter.
*
* @param conf a configuration
* @return a filter for the unbound record filter specified in conf
*/
public static Filter getFilter(ParquetConfiguration conf) {
return FilterCompat.get(getFilterPredicate(conf), getUnboundRecordFilterInstance(conf));
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
/*
* 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.conf;

/**
* Stores the read configuration property keys used to configure how Parquet
* files are read. This class contains only constants and is decoupled from the
* legacy Hadoop {@code Configuration} and {@code InputFormat} machinery, so that
* the configuration properties can be used without loading {@code FileInputFormat}
* and its transitive dependencies.
*/
public final class ParquetInputProperties {

/** key to configure the ReadSupport implementation */
public static final String READ_SUPPORT_CLASS = "parquet.read.support.class";

/** key to configure the filter */
public static final String UNBOUND_RECORD_FILTER = "parquet.read.filter";

/** key to configure type checking for conflicting schemas (default: true) */
public static final String STRICT_TYPE_CHECKING = "parquet.strict.typing";

/** key to configure the filter predicate */
public static final String FILTER_PREDICATE = "parquet.private.read.filter.predicate";

/** key to configure whether record-level filtering is enabled */
public static final String RECORD_FILTERING_ENABLED = "parquet.filter.record-level.enabled";

/** key to configure whether row group stats filtering is enabled */
public static final String STATS_FILTERING_ENABLED = "parquet.filter.stats.enabled";

/** key to configure whether row group dictionary filtering is enabled */
public static final String DICTIONARY_FILTERING_ENABLED = "parquet.filter.dictionary.enabled";

/** key to configure whether column index filtering of pages is enabled */
public static final String COLUMN_INDEX_FILTERING_ENABLED = "parquet.filter.columnindex.enabled";

/** key to configure whether page level checksum verification is enabled */
public static final String PAGE_VERIFY_CHECKSUM_ENABLED = "parquet.page.verify-checksum.enabled";

/** key to configure whether row group bloom filtering is enabled */
public static final String BLOOM_FILTERING_ENABLED = "parquet.filter.bloom.enabled";

/** Key to configure if off-heap buffer should be used for decryption */
public static final String OFF_HEAP_DECRYPT_BUFFER_ENABLED = "parquet.decrypt.off-heap.buffer.enabled";

/** Key to enable/disable vectored io while reading parquet files: {@value}. */
public static final String HADOOP_VECTORED_IO_ENABLED = "parquet.hadoop.vectored.io.enabled";

/** Default value of parquet.hadoop.vectored.io.enabled is {@value}. */
public static final boolean HADOOP_VECTORED_IO_DEFAULT = true;

private ParquetInputProperties() {}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@
package org.apache.parquet.hadoop;

import static java.lang.String.format;
import static org.apache.parquet.hadoop.ParquetInputFormat.RECORD_FILTERING_ENABLED;
import static org.apache.parquet.hadoop.ParquetInputFormat.STRICT_TYPE_CHECKING;
import static org.apache.parquet.conf.ParquetInputProperties.RECORD_FILTERING_ENABLED;
import static org.apache.parquet.conf.ParquetInputProperties.STRICT_TYPE_CHECKING;

import java.io.IOException;
import java.util.Collections;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,8 +51,9 @@
import org.apache.parquet.Preconditions;
import org.apache.parquet.conf.HadoopParquetConfiguration;
import org.apache.parquet.conf.ParquetConfiguration;
import org.apache.parquet.conf.ParquetInputFilters;
import org.apache.parquet.conf.ParquetInputProperties;
import org.apache.parquet.filter.UnboundRecordFilter;
import org.apache.parquet.filter2.compat.FilterCompat;
import org.apache.parquet.filter2.compat.FilterCompat.Filter;
import org.apache.parquet.filter2.compat.RowGroupFilter;
import org.apache.parquet.filter2.predicate.FilterPredicate;
Expand Down Expand Up @@ -82,10 +83,10 @@
* It must be a subset of the original schema. Only the columns needed to reconstruct the records with the requestedSchema will be scanned.
*
* @param <T> the type of the materialized records
* @see #READ_SUPPORT_CLASS
* @see #UNBOUND_RECORD_FILTER
* @see #STRICT_TYPE_CHECKING
* @see #FILTER_PREDICATE
* @see ParquetInputProperties#READ_SUPPORT_CLASS
* @see ParquetInputProperties#UNBOUND_RECORD_FILTER
* @see ParquetInputProperties#STRICT_TYPE_CHECKING
* @see ParquetInputProperties#FILTER_PREDICATE
* @see #TASK_SIDE_METADATA
*/
public class ParquetInputFormat<T> extends FileInputFormat<Void, T> {
Expand All @@ -94,58 +95,80 @@ public class ParquetInputFormat<T> extends FileInputFormat<Void, T> {

/**
* key to configure the ReadSupport implementation
* @deprecated use {@link ParquetInputProperties#READ_SUPPORT_CLASS}
*/
public static final String READ_SUPPORT_CLASS = "parquet.read.support.class";
@Deprecated
public static final String READ_SUPPORT_CLASS = ParquetInputProperties.READ_SUPPORT_CLASS;

/**
* key to configure the filter
* @deprecated use {@link ParquetInputProperties#UNBOUND_RECORD_FILTER}
*/
public static final String UNBOUND_RECORD_FILTER = "parquet.read.filter";
@Deprecated
public static final String UNBOUND_RECORD_FILTER = ParquetInputProperties.UNBOUND_RECORD_FILTER;

/**
* key to configure type checking for conflicting schemas (default: true)
* @deprecated use {@link ParquetInputProperties#STRICT_TYPE_CHECKING}
*/
public static final String STRICT_TYPE_CHECKING = "parquet.strict.typing";
@Deprecated
public static final String STRICT_TYPE_CHECKING = ParquetInputProperties.STRICT_TYPE_CHECKING;

/**
* key to configure the filter predicate
* @deprecated use {@link ParquetInputProperties#FILTER_PREDICATE}
*/
public static final String FILTER_PREDICATE = "parquet.private.read.filter.predicate";
@Deprecated
public static final String FILTER_PREDICATE = ParquetInputProperties.FILTER_PREDICATE;

/**
* key to configure whether record-level filtering is enabled
* @deprecated use {@link ParquetInputProperties#RECORD_FILTERING_ENABLED}
*/
public static final String RECORD_FILTERING_ENABLED = "parquet.filter.record-level.enabled";
@Deprecated
public static final String RECORD_FILTERING_ENABLED = ParquetInputProperties.RECORD_FILTERING_ENABLED;

/**
* key to configure whether row group stats filtering is enabled
* @deprecated use {@link ParquetInputProperties#STATS_FILTERING_ENABLED}
*/
public static final String STATS_FILTERING_ENABLED = "parquet.filter.stats.enabled";
@Deprecated
public static final String STATS_FILTERING_ENABLED = ParquetInputProperties.STATS_FILTERING_ENABLED;

/**
* key to configure whether row group dictionary filtering is enabled
* @deprecated use {@link ParquetInputProperties#DICTIONARY_FILTERING_ENABLED}
*/
public static final String DICTIONARY_FILTERING_ENABLED = "parquet.filter.dictionary.enabled";
@Deprecated
public static final String DICTIONARY_FILTERING_ENABLED = ParquetInputProperties.DICTIONARY_FILTERING_ENABLED;

/**
* key to configure whether column index filtering of pages is enabled
* @deprecated use {@link ParquetInputProperties#COLUMN_INDEX_FILTERING_ENABLED}
*/
public static final String COLUMN_INDEX_FILTERING_ENABLED = "parquet.filter.columnindex.enabled";
@Deprecated
public static final String COLUMN_INDEX_FILTERING_ENABLED = ParquetInputProperties.COLUMN_INDEX_FILTERING_ENABLED;

/**
* key to configure whether page level checksum verification is enabled
* @deprecated use {@link ParquetInputProperties#PAGE_VERIFY_CHECKSUM_ENABLED}
*/
public static final String PAGE_VERIFY_CHECKSUM_ENABLED = "parquet.page.verify-checksum.enabled";
@Deprecated
public static final String PAGE_VERIFY_CHECKSUM_ENABLED = ParquetInputProperties.PAGE_VERIFY_CHECKSUM_ENABLED;

/**
* key to configure whether row group bloom filtering is enabled
* @deprecated use {@link ParquetInputProperties#BLOOM_FILTERING_ENABLED}
*/
public static final String BLOOM_FILTERING_ENABLED = "parquet.filter.bloom.enabled";
@Deprecated
public static final String BLOOM_FILTERING_ENABLED = ParquetInputProperties.BLOOM_FILTERING_ENABLED;

/**
* Key to configure if off-heap buffer should be used for decryption
* @deprecated use {@link ParquetInputProperties#OFF_HEAP_DECRYPT_BUFFER_ENABLED}
*/
public static final String OFF_HEAP_DECRYPT_BUFFER_ENABLED = "parquet.decrypt.off-heap.buffer.enabled";
@Deprecated
public static final String OFF_HEAP_DECRYPT_BUFFER_ENABLED = ParquetInputProperties.OFF_HEAP_DECRYPT_BUFFER_ENABLED;

/**
* key to turn on or off task side metadata loading (default true)
Expand All @@ -164,13 +187,17 @@ public class ParquetInputFormat<T> extends FileInputFormat<Void, T> {
/**
* Key to enable/disable vectored io while reading parquet files:
* {@value}.
* @deprecated use {@link ParquetInputProperties#HADOOP_VECTORED_IO_ENABLED}
*/
public static final String HADOOP_VECTORED_IO_ENABLED = "parquet.hadoop.vectored.io.enabled";
@Deprecated
public static final String HADOOP_VECTORED_IO_ENABLED = ParquetInputProperties.HADOOP_VECTORED_IO_ENABLED;

/**
* Default value of parquet.hadoop.vectored.io.enabled is {@value}.
* @deprecated use {@link ParquetInputProperties#HADOOP_VECTORED_IO_DEFAULT}
*/
public static final boolean HADOOP_VECTORED_IO_DEFAULT = true;
@Deprecated
public static final boolean HADOOP_VECTORED_IO_DEFAULT = ParquetInputProperties.HADOOP_VECTORED_IO_DEFAULT;

public static void setTaskSideMetaData(Job job, boolean taskSideMetadata) {
ContextUtil.getConfiguration(job).setBoolean(TASK_SIDE_METADATA, taskSideMetadata);
Expand Down Expand Up @@ -200,20 +227,7 @@ public static void setUnboundRecordFilter(Job job, Class<? extends UnboundRecord
*/
@Deprecated
public static Class<?> getUnboundRecordFilter(Configuration configuration) {
return ConfigurationUtil.getClassFromConfig(configuration, UNBOUND_RECORD_FILTER, UnboundRecordFilter.class);
}

private static UnboundRecordFilter getUnboundRecordFilterInstance(ParquetConfiguration configuration) {
Class<?> clazz =
ConfigurationUtil.getClassFromConfig(configuration, UNBOUND_RECORD_FILTER, UnboundRecordFilter.class);
if (clazz == null) {
return null;
}
try {
return (UnboundRecordFilter) clazz.newInstance();
} catch (InstantiationException | IllegalAccessException e) {
throw new BadConfigurationException("could not instantiate unbound record filter class", e);
}
return ParquetInputFilters.getUnboundRecordFilter(new HadoopParquetConfiguration(configuration));
}

public static void setReadSupportClass(JobConf conf, Class<?> readSupportClass) {
Expand All @@ -238,15 +252,7 @@ public static void setFilterPredicate(Configuration configuration, FilterPredica
}

private static FilterPredicate getFilterPredicate(Configuration configuration) {
return getFilterPredicate(new HadoopParquetConfiguration(configuration));
}

private static FilterPredicate getFilterPredicate(ParquetConfiguration configuration) {
try {
return SerializationUtil.readObjectFromConfAsBase64(FILTER_PREDICATE, configuration);
} catch (IOException e) {
throw new RuntimeException(e);
}
return ParquetInputFilters.getFilterPredicate(new HadoopParquetConfiguration(configuration));
}

/**
Expand All @@ -260,8 +266,17 @@ public static Filter getFilter(Configuration conf) {
return getFilter(new HadoopParquetConfiguration(conf));
}

/**
* Returns a non-null Filter, which is a wrapper around either a
* FilterPredicate, an UnboundRecordFilter, or a no-op filter.
*
* @param conf a configuration
* @return a filter for the unbound record filter specified in conf
* @deprecated use {@link ParquetInputFilters#getFilter(ParquetConfiguration)}
*/
@Deprecated
public static Filter getFilter(ParquetConfiguration conf) {
return FilterCompat.get(getFilterPredicate(conf), getUnboundRecordFilterInstance(conf));
return ParquetInputFilters.getFilter(conf);
}

private LruCache<FileStatusWrapper, FootersCacheValue> footersCache;
Expand Down
Loading