From c5ee65d33f79b6194227b57eb5cd8a79f0d7ed62 Mon Sep 17 00:00:00 2001 From: Costas Zarifis Date: Mon, 28 Sep 2026 06:49:35 +0000 Subject: [PATCH 1/2] GH-3826: Add Bloom path collision reproduction --- .../parquet/hadoop/TestParquetWriter.java | 110 ++++++++++++++++++ .../bloom-filter-path-collision.json | 6 + 2 files changed, 116 insertions(+) create mode 100644 parquet-hadoop/src/test/resources/bloom-filter-path-collision.json diff --git a/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestParquetWriter.java b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestParquetWriter.java index c501298e76..b4a7fc9b94 100644 --- a/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestParquetWriter.java +++ b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestParquetWriter.java @@ -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; @@ -335,6 +337,114 @@ public void testParquetFileWithBloomFilter() throws IOException { } } + @Test + public void testBloomFiltersKeepCollidingDotStringPathsDistinct() throws IOException { + Map 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 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 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 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 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; diff --git a/parquet-hadoop/src/test/resources/bloom-filter-path-collision.json b/parquet-hadoop/src/test/resources/bloom-filter-path-collision.json new file mode 100644 index 0000000000..0d6a55d7c0 --- /dev/null +++ b/parquet-hadoop/src/test/resources/bloom-filter-path-collision.json @@ -0,0 +1,6 @@ +{ + "topLevelPath" : [ "a.b" ], + "topLevelContainsOwnBloomValues" : true, + "nestedPath" : [ "a", "b" ], + "nestedContainsOwnBloomValues" : true +} From 19217b47407f94a31247c248124484e1627c8989 Mon Sep 17 00:00:00 2001 From: Costas Zarifis Date: Mon, 28 Sep 2026 07:24:42 +0000 Subject: [PATCH 2/2] GH-3826: Preserve component paths for Bloom filters --- .../parquet/hadoop/ParquetFileWriter.java | 24 +++++++++++++------ .../hadoop/rewrite/ParquetRewriter.java | 2 +- 2 files changed, 18 insertions(+), 8 deletions(-) diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java index 82f4577b83..00034696d4 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java @@ -134,7 +134,7 @@ public static enum Mode { private final List> offsetIndexes = new ArrayList<>(); // The Bloom filters - private final List> bloomFilters = new ArrayList<>(); + private final List> bloomFilters = new ArrayList<>(); // The file encryptor private final InternalFileEncryptor fileEncryptor; @@ -147,7 +147,7 @@ public static enum Mode { private List currentOffsetIndexes; // The Bloom filter for the actual block - private Map currentBloomFilters; + private Map currentBloomFilters; // row group data set at the start of a row group private long currentRecordCount; // set in startBlock @@ -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); } @@ -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.", @@ -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); @@ -2036,7 +2046,7 @@ private static void serializeOffsetIndexes( } private static void serializeBloomFilters( - List> bloomFilters, + List> bloomFilters, List blocks, PositionOutputStream out, InternalFileEncryptor fileEncryptor) @@ -2045,11 +2055,11 @@ private static void serializeBloomFilters( for (int bIndex = 0, bSize = blocks.size(); bIndex < bSize; ++bIndex) { BlockMetaData block = blocks.get(bIndex); List columns = block.getColumns(); - Map blockBloomFilters = bloomFilters.get(bIndex); + Map 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; } diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/rewrite/ParquetRewriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/rewrite/ParquetRewriter.java index 88e41626dd..e595dabef4 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/rewrite/ParquetRewriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/rewrite/ParquetRewriter.java @@ -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());