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
@@ -0,0 +1,124 @@
/*
* 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.iceberg.mr.hive;

import java.io.IOException;
import java.util.Collection;
import java.util.List;
import org.apache.hadoop.fs.Path;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.FileFormat;
import org.apache.iceberg.Schema;
import org.apache.iceberg.Table;
import org.apache.iceberg.TableProperties;
import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.expressions.Expressions;
import org.apache.iceberg.mr.hive.test.TestTables.TestTableType;
import org.apache.iceberg.parquet.ParquetBloomRowGroupFilter;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.types.Types;
import org.apache.parquet.column.values.bloomfilter.BloomFilter;
import org.apache.parquet.hadoop.BloomFilterReader;
import org.apache.parquet.hadoop.ParquetFileReader;
import org.apache.parquet.hadoop.metadata.BlockMetaData;
import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData;
import org.apache.parquet.hadoop.util.HadoopInputFile;
import org.apache.parquet.schema.MessageType;
import org.junit.Assert;
import org.junit.Test;
import org.junit.runners.Parameterized.Parameters;

import static org.apache.iceberg.types.Types.NestedField.optional;
import static org.apache.iceberg.types.Types.NestedField.required;

/**
* Verifies that Hive inserts into Parquet Iceberg tables honor the
* {@code write.parquet.bloom-filter-enabled.column.*} table properties: the written files must contain working
* bloom filters that Iceberg's read-side row group filter can prune on.
*/
public class TestHiveIcebergParquetBloomFilter extends HiveIcebergStorageHandlerWithEngineBase {

private static final long PRESENT_ID = 42L;
private static final long ABSENT_ID = 12345678L;

@Parameters(name = "fileFormat={0}, catalog={1}, isVectorized={2}, formatVersion={3}")
public static Collection<Object[]> parameters() {
return HiveIcebergStorageHandlerWithEngineBase.getParameters(p ->
p.fileFormat() == FileFormat.PARQUET && p.testTableType() == TestTableType.HIVE_CATALOG &&
p.formatVersion() == 2);
}

@Test
public void testBloomFilterWrittenByHiveInsert() throws IOException {
Schema schema = new Schema(
required(1, "id", Types.LongType.get()),
optional(2, "name", Types.StringType.get()));

testTables.createTable(shell, "bloom_test", schema, fileFormat, ImmutableList.of(), formatVersion,
ImmutableMap.of(
TableProperties.PARQUET_BLOOM_FILTER_COLUMN_ENABLED_PREFIX + "id", "true",
TableProperties.PARQUET_BLOOM_FILTER_COLUMN_FPP_PREFIX + "id", "0.01"));

shell.executeStatement("INSERT INTO bloom_test VALUES (1, 'a'), (" + PRESENT_ID + ", 'b'), (100, 'c')");

Table table = testTables.loadTable(TableIdentifier.of("default", "bloom_test"));
List<DataFile> dataFiles = Lists.newArrayList(table.currentSnapshot().addedDataFiles(table.io()));
Assert.assertEquals(1, dataFiles.size());

HadoopInputFile inputFile = HadoopInputFile.fromPath(new Path(dataFiles.get(0).location()), shell.getHiveConf());
try (ParquetFileReader reader = ParquetFileReader.open(inputFile)) {
MessageType fileSchema = reader.getFooter().getFileMetaData().getSchema();
List<BlockMetaData> rowGroups = reader.getFooter().getBlocks();
Assert.assertFalse(rowGroups.isEmpty());

for (BlockMetaData rowGroup : rowGroups) {
BloomFilterReader bloomReader = reader.getBloomFilterDataReader(rowGroup);

BloomFilter bloom = bloomReader.readBloomFilter(columnChunk(rowGroup, "id"));
Assert.assertNotNull("Bloom filter should be written for the enabled column", bloom);
Assert.assertTrue(bloom.findHash(bloom.hash(PRESENT_ID)));
Assert.assertFalse(bloom.findHash(bloom.hash(ABSENT_ID)));

Assert.assertNull("Bloom filter should not be written for a column where it was not enabled",
bloomReader.readBloomFilter(columnChunk(rowGroup, "name")));

Assert.assertTrue(new ParquetBloomRowGroupFilter(schema, Expressions.equal("id", PRESENT_ID))
.shouldRead(fileSchema, rowGroup, bloomReader));
Assert.assertFalse("Row group should be prunable for a value not in the bloom filter",
new ParquetBloomRowGroupFilter(schema, Expressions.equal("id", ABSENT_ID))
.shouldRead(fileSchema, rowGroup, bloomReader));
}
}

List<Object[]> rows = shell.executeStatement("SELECT name FROM bloom_test WHERE id = " + PRESENT_ID);
Assert.assertEquals(1, rows.size());
Assert.assertEquals("b", rows.get(0)[0]);
Assert.assertTrue(shell.executeStatement("SELECT * FROM bloom_test WHERE id = " + ABSENT_ID).isEmpty());
}

private static ColumnChunkMetaData columnChunk(BlockMetaData rowGroup, String columnName) {
return rowGroup.getColumns().stream()
.filter(column -> column.getPath().toDotString().equals(columnName))
.findAny()
.orElseThrow();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
-- Mask random uuid
--! qt:replace:/(\s+'uuid'=')\S+('\s*)/$1#Masked#$2/

-- Parquet bloom filter write properties on Iceberg tables: verifies the properties are accepted,
-- survive in HMS, and inserts/point lookups work with bloom filters enabled.
-- Bloom filter presence in the data files is asserted by TestHiveIcebergParquetBloomFilter.

drop table if exists tbl_bloom;
create external table tbl_bloom(id bigint, name string) stored by iceberg stored as parquet
tblproperties ('format-version'='2',
'write.parquet.bloom-filter-enabled.column.id'='true',
'write.parquet.bloom-filter-fpp.column.id'='0.05');

show create table tbl_bloom;

insert into tbl_bloom values (1, 'one'), (42, 'answer'), (100, 'hundred'), (12345678, 'big');

select name from tbl_bloom where id = 42;
select count(*) from tbl_bloom where id = 43;
select * from tbl_bloom order by id;

-- enable bloom filter on another column, subsequent writes pick it up
alter table tbl_bloom set tblproperties ('write.parquet.bloom-filter-enabled.column.name'='true');
insert into tbl_bloom values (200, 'two hundred');

select id from tbl_bloom where name = 'two hundred';
select count(*) from tbl_bloom;

drop table tbl_bloom;
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
-- Parquet bloom filter pruning under vectorized LLAP execution, where the filters are served from the
-- LLAP metadata cache after the first read of a file.
set hive.llap.io.enabled=true;
set hive.vectorized.execution.enabled=true;

DROP TABLE IF EXISTS llap_bloom_parquet PURGE;

CREATE EXTERNAL TABLE llap_bloom_parquet (id bigint, name string)
STORED BY ICEBERG STORED AS PARQUET
TBLPROPERTIES ('format-version'='2', 'write.parquet.bloom-filter-enabled.column.id'='true');

INSERT INTO llap_bloom_parquet VALUES
(2, 'two'), (4, 'four'), (6, 'six'), (8, 'eight'), (10, 'ten');

-- absent from the bloom filter but inside the min/max range, so only the bloom filter can prune it;
-- this first read fills the cache
SELECT count(*) FROM llap_bloom_parquet WHERE id = 5;

-- under cache.only the reader may not fall back to the file, so this answers only if the bloom filter
-- itself came from the cache
set hive.llap.io.cache.only=true;
SELECT count(*) FROM llap_bloom_parquet WHERE id = 5;
set hive.llap.io.cache.only=false;

-- a value the bloom filter does contain must survive pruning; count(*) keeps this off the fetch-task
-- path, which runs in the client JVM and would never reach the LLAP reader
SELECT count(*) FROM llap_bloom_parquet WHERE id = 6;

DROP TABLE llap_bloom_parquet PURGE;

-- A file of several row groups, where statistics leave a different row group standing per predicate. Each
-- filter is cached under its own offset, so serving one never depends on which query cached it.
DROP TABLE IF EXISTS llap_bloom_multi PURGE;

CREATE EXTERNAL TABLE llap_bloom_multi (id bigint, name string)
STORED BY ICEBERG STORED AS PARQUET
TBLPROPERTIES ('format-version'='2', 'write.parquet.bloom-filter-enabled.column.id'='true',
'write.parquet.bloom-filter-max-bytes'='1024', 'write.parquet.row-group-size-bytes'='1024');

INSERT INTO llap_bloom_multi
SELECT pos * 2, concat('n', pos) FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val;

-- odd ids are absent everywhere, and each lands in a different row group
SELECT count(*) FROM llap_bloom_multi WHERE id = 51;
SELECT count(*) FROM llap_bloom_multi WHERE id = 651;

set hive.llap.io.cache.only=true;
SELECT count(*) FROM llap_bloom_multi WHERE id = 51;
SELECT count(*) FROM llap_bloom_multi WHERE id = 651;
set hive.llap.io.cache.only=false;

-- even ids are present, and must survive pruning against filters served from the cache
SELECT count(*) FROM llap_bloom_multi WHERE id = 4;
SELECT count(*) FROM llap_bloom_multi WHERE id = 700;

DROP TABLE llap_bloom_multi PURGE;
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
PREHOOK: query: drop table if exists tbl_bloom
PREHOOK: type: DROPTABLE
PREHOOK: Output: database:default
POSTHOOK: query: drop table if exists tbl_bloom
POSTHOOK: type: DROPTABLE
POSTHOOK: Output: database:default
PREHOOK: query: create external table tbl_bloom(id bigint, name string) stored by iceberg stored as parquet
tblproperties ('format-version'='2',
'write.parquet.bloom-filter-enabled.column.id'='true',
'write.parquet.bloom-filter-fpp.column.id'='0.05')
PREHOOK: type: CREATETABLE
PREHOOK: Output: database:default
PREHOOK: Output: default@tbl_bloom
POSTHOOK: query: create external table tbl_bloom(id bigint, name string) stored by iceberg stored as parquet
tblproperties ('format-version'='2',
'write.parquet.bloom-filter-enabled.column.id'='true',
'write.parquet.bloom-filter-fpp.column.id'='0.05')
POSTHOOK: type: CREATETABLE
POSTHOOK: Output: database:default
POSTHOOK: Output: default@tbl_bloom
PREHOOK: query: show create table tbl_bloom
PREHOOK: type: SHOW_CREATETABLE
PREHOOK: Input: default@tbl_bloom
POSTHOOK: query: show create table tbl_bloom
POSTHOOK: type: SHOW_CREATETABLE
POSTHOOK: Input: default@tbl_bloom
CREATE EXTERNAL TABLE `tbl_bloom`(
`id` bigint,
`name` string)
ROW FORMAT SERDE
'org.apache.iceberg.mr.hive.HiveIcebergSerDe'
STORED BY
'org.apache.iceberg.mr.hive.HiveIcebergStorageHandler'

LOCATION
'hdfs://### HDFS PATH ###'
TBLPROPERTIES (
'bucketing_version'='2',
'current-schema'='{"type":"struct","schema-id":0,"fields":[{"id":1,"name":"id","required":false,"type":"long"},{"id":2,"name":"name","required":false,"type":"string"}]}',
'format-version'='2',
'metadata_location'='hdfs://### HDFS PATH ###',
'parquet.compression'='zstd',
'serialization.format'='1',
'snapshot-count'='0',
'table_type'='ICEBERG',
#### A masked pattern was here ####
'uuid'='#Masked#',
'write.delete.mode'='merge-on-read',
'write.format.default'='parquet',
'write.merge.mode'='merge-on-read',
'write.metadata.delete-after-commit.enabled'='true',
'write.parquet.bloom-filter-enabled.column.id'='true',
'write.parquet.bloom-filter-fpp.column.id'='0.05',
'write.update.mode'='merge-on-read')
PREHOOK: query: insert into tbl_bloom values (1, 'one'), (42, 'answer'), (100, 'hundred'), (12345678, 'big')
PREHOOK: type: QUERY
PREHOOK: Input: _dummy_database@_dummy_table
PREHOOK: Output: default@tbl_bloom
POSTHOOK: query: insert into tbl_bloom values (1, 'one'), (42, 'answer'), (100, 'hundred'), (12345678, 'big')
POSTHOOK: type: QUERY
POSTHOOK: Input: _dummy_database@_dummy_table
POSTHOOK: Output: default@tbl_bloom
PREHOOK: query: select name from tbl_bloom where id = 42
PREHOOK: type: QUERY
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: select name from tbl_bloom where id = 42
POSTHOOK: type: QUERY
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: hdfs://### HDFS PATH ###
answer
PREHOOK: query: select count(*) from tbl_bloom where id = 43
PREHOOK: type: QUERY
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: select count(*) from tbl_bloom where id = 43
POSTHOOK: type: QUERY
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: hdfs://### HDFS PATH ###
0
PREHOOK: query: select * from tbl_bloom order by id
PREHOOK: type: QUERY
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: select * from tbl_bloom order by id
POSTHOOK: type: QUERY
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: hdfs://### HDFS PATH ###
1 one
42 answer
100 hundred
12345678 big
PREHOOK: query: alter table tbl_bloom set tblproperties ('write.parquet.bloom-filter-enabled.column.name'='true')
PREHOOK: type: ALTERTABLE_PROPERTIES
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: default@tbl_bloom
POSTHOOK: query: alter table tbl_bloom set tblproperties ('write.parquet.bloom-filter-enabled.column.name'='true')
POSTHOOK: type: ALTERTABLE_PROPERTIES
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: default@tbl_bloom
PREHOOK: query: insert into tbl_bloom values (200, 'two hundred')
PREHOOK: type: QUERY
PREHOOK: Input: _dummy_database@_dummy_table
PREHOOK: Output: default@tbl_bloom
POSTHOOK: query: insert into tbl_bloom values (200, 'two hundred')
POSTHOOK: type: QUERY
POSTHOOK: Input: _dummy_database@_dummy_table
POSTHOOK: Output: default@tbl_bloom
PREHOOK: query: select id from tbl_bloom where name = 'two hundred'
PREHOOK: type: QUERY
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: select id from tbl_bloom where name = 'two hundred'
POSTHOOK: type: QUERY
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: hdfs://### HDFS PATH ###
200
PREHOOK: query: select count(*) from tbl_bloom
PREHOOK: type: QUERY
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: hdfs://### HDFS PATH ###
POSTHOOK: query: select count(*) from tbl_bloom
POSTHOOK: type: QUERY
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: hdfs://### HDFS PATH ###
5
PREHOOK: query: drop table tbl_bloom
PREHOOK: type: DROPTABLE
PREHOOK: Input: default@tbl_bloom
PREHOOK: Output: database:default
PREHOOK: Output: default@tbl_bloom
POSTHOOK: query: drop table tbl_bloom
POSTHOOK: type: DROPTABLE
POSTHOOK: Input: default@tbl_bloom
POSTHOOK: Output: database:default
POSTHOOK: Output: default@tbl_bloom
Loading
Loading