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 @@ -134,7 +134,7 @@ public static enum Mode {
private final List<List<OffsetIndex>> offsetIndexes = new ArrayList<>();

// The Bloom filters
private final List<Map<String, BloomFilter>> bloomFilters = new ArrayList<>();
private final List<Map<ColumnPath, BloomFilter>> bloomFilters = new ArrayList<>();

// The file encryptor
private final InternalFileEncryptor fileEncryptor;
Expand All @@ -147,7 +147,7 @@ public static enum Mode {
private List<OffsetIndex> currentOffsetIndexes;

// The Bloom filter for the actual block
private Map<String, BloomFilter> currentBloomFilters;
private Map<ColumnPath, BloomFilter> currentBloomFilters;

// row group data set at the start of a row group
private long currentRecordCount; // set in startBlock
Expand Down Expand Up @@ -1078,6 +1078,16 @@ public void writeDataPage(
* @param bloomFilter the bloom filter of column values
*/
public void addBloomFilter(String column, BloomFilter bloomFilter) {
addBloomFilterForPath(ColumnPath.fromDotString(column), bloomFilter);
}

/**
* Add a Bloom filter that will be written out.
*
* @param column the component-based column path
* @param bloomFilter the bloom filter of column values
*/
public void addBloomFilterForPath(ColumnPath column, BloomFilter bloomFilter) {
currentBloomFilters.put(column, bloomFilter);
}

Expand Down Expand Up @@ -1531,7 +1541,7 @@ void writeColumnChunk(
}
}
if (isWriteBloomFilter) {
currentBloomFilters.put(String.join(".", descriptor.getPath()), bloomFilter);
currentBloomFilters.put(ColumnPath.get(descriptor.getPath()), bloomFilter);
} else {
LOG.info(
"No need to write bloom filter because column {} data pages are all encoded as dictionary.",
Expand Down Expand Up @@ -1813,7 +1823,7 @@ public void appendColumnChunk(

copy(from, out, start, length);

currentBloomFilters.put(String.join(".", descriptor.getPath()), bloomFilter);
currentBloomFilters.put(ColumnPath.get(descriptor.getPath()), bloomFilter);
currentColumnIndexes.add(columnIndex);
currentOffsetIndexes.add(effectiveOffsetIndex);

Expand Down Expand Up @@ -2036,7 +2046,7 @@ private static void serializeOffsetIndexes(
}

private static void serializeBloomFilters(
List<Map<String, BloomFilter>> bloomFilters,
List<Map<ColumnPath, BloomFilter>> bloomFilters,
List<BlockMetaData> blocks,
PositionOutputStream out,
InternalFileEncryptor fileEncryptor)
Expand All @@ -2045,11 +2055,11 @@ private static void serializeBloomFilters(
for (int bIndex = 0, bSize = blocks.size(); bIndex < bSize; ++bIndex) {
BlockMetaData block = blocks.get(bIndex);
List<ColumnChunkMetaData> columns = block.getColumns();
Map<String, BloomFilter> blockBloomFilters = bloomFilters.get(bIndex);
Map<ColumnPath, BloomFilter> blockBloomFilters = bloomFilters.get(bIndex);
if (blockBloomFilters.isEmpty()) continue;
for (int cIndex = 0, cSize = columns.size(); cIndex < cSize; ++cIndex) {
ColumnChunkMetaData column = columns.get(cIndex);
BloomFilter bloomFilter = blockBloomFilters.get(column.getPath().toDotString());
BloomFilter bloomFilter = blockBloomFilters.get(column.getPath());
if (bloomFilter == null) {
continue;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -614,7 +614,7 @@ private void processChunk(
}

if (bloomFilter != null) {
writer.addBloomFilter(normalizeFieldsInPath(chunk.getPath()).toDotString(), bloomFilter);
writer.addBloomFilterForPath(normalizeFieldsInPath(chunk.getPath()), bloomFilter);
}

reader.setStreamPosition(chunk.getStartingPos());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,9 @@

import com.google.common.collect.ImmutableMap;
import java.io.IOException;
import java.io.InputStream;
import java.lang.reflect.Field;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
Expand Down Expand Up @@ -335,6 +337,114 @@ public void testParquetFileWithBloomFilter() throws IOException {
}
}

@Test
public void testBloomFiltersKeepCollidingDotStringPathsDistinct() throws IOException {
Map<ColumnPath, BloomFilter> bloomFilters = writeBloomFiltersForCollidingPaths();
BloomFilter topLevelBloom = bloomFilters.get(ColumnPath.get("a.b"));
BloomFilter nestedBloom = bloomFilters.get(ColumnPath.get("a", "b"));

for (int valueId = 0; valueId < 10; valueId++) {
long topLevelHash = LongHashFunction.xx(0)
.hashBytes(Binary.fromString("top-" + valueId).toByteBuffer());
long nestedHash = LongHashFunction.xx(0)
.hashBytes(Binary.fromString("nested-" + valueId).toByteBuffer());
assertThat(topLevelBloom.findHash(topLevelHash)).isTrue();
assertThat(nestedBloom.findHash(nestedHash)).isTrue();
}
}

@Test
public void testBloomFilterPathCollisionGolden() throws IOException {
Map<ColumnPath, BloomFilter> bloomFilters = writeBloomFiltersForCollidingPaths();
String actual = String.format(
"{\n"
+ " \"topLevelPath\" : [ \"a.b\" ],\n"
+ " \"topLevelContainsOwnBloomValues\" : %s,\n"
+ " \"nestedPath\" : [ \"a\", \"b\" ],\n"
+ " \"nestedContainsOwnBloomValues\" : %s\n"
+ "}\n",
containsAllBloomValues(bloomFilters.get(ColumnPath.get("a.b")), "top-"),
containsAllBloomValues(bloomFilters.get(ColumnPath.get("a", "b")), "nested-"));

try (InputStream input = TestParquetWriter.class.getResourceAsStream("/bloom-filter-path-collision.json")) {
assertThat(input).isNotNull();
String expected = new String(input.readAllBytes(), StandardCharsets.UTF_8);
assertThat(actual).isEqualTo(expected);
}
}

private Map<ColumnPath, BloomFilter> writeBloomFiltersForCollidingPaths() throws IOException {
// Parquet identifies columns by path components. These two paths are different:
//
// top-level field named "a.b" -> ["a.b"]
// nested field "b" in "a" -> ["a", "b"]
//
// Flattening either path with '.' produces the same string, "a.b". When Bloom filters are
// stored in a map keyed by that flattened string, the second column overwrites the first
// column's filter and both footer entries can be associated with the same filter.
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);
GroupFactory factory = new SimpleGroupFactory(schema);

Path path = newTempPath();
// Disable dictionary encoding so the writer emits a Bloom filter for each column. Use
// disjoint value prefixes so the test can detect when one column receives the other's filter.
try (ParquetWriter<Group> writer = ExampleParquetWriter.builder(path)
.withAllocator(allocator)
.withConf(conf)
.withDictionaryEncoding(false)
.withBloomFilterEnabled(true)
.build()) {
for (int i = 0; i < 100; i++) {
int valueId = i % 10;
Group group = factory.newGroup().append("a.b", "top-" + valueId);
group.addGroup("a").append("b", "nested-" + valueId);
writer.write(group);
}
}

try (ParquetFileReader reader = ParquetFileReader.open(HadoopInputFile.fromPath(path, conf))) {
BlockMetaData block = reader.getFooter().getBlocks().get(0);
// Locate columns by their component-based ColumnPath. Using toDotString() here would repeat
// the bug and make both columns indistinguishable to the test.
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();
BloomFilter topLevelBloom = reader.readBloomFilter(topLevelColumn);
BloomFilter nestedBloom = reader.readBloomFilter(nestedColumn);
Map<ColumnPath, BloomFilter> bloomFilters = new HashMap<>();
bloomFilters.put(ColumnPath.get("a.b"), topLevelBloom);
bloomFilters.put(ColumnPath.get("a", "b"), nestedBloom);
return bloomFilters;
}
}

private static boolean containsAllBloomValues(BloomFilter bloomFilter, String prefix) {
for (int valueId = 0; valueId < 10; valueId++) {
long hash = LongHashFunction.xx(0)
.hashBytes(Binary.fromString(prefix + valueId).toByteBuffer());
if (!bloomFilter.findHash(hash)) {
return false;
}
}
return true;
}

@Test
public void testParquetFileWithBloomFilterWithFpp() throws IOException {
int buildBloomFilterCount = 100000;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
{
"topLevelPath" : [ "a.b" ],
"topLevelContainsOwnBloomValues" : true,
"nestedPath" : [ "a", "b" ],
"nestedContainsOwnBloomValues" : true
}