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 @@ -31,6 +31,7 @@
import org.apache.parquet.column.page.DictionaryPageReadStore;
import org.apache.parquet.hadoop.metadata.BlockMetaData;
import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData;
import org.apache.parquet.hadoop.metadata.ColumnPath;
import org.apache.parquet.io.ParquetDecodingException;

/**
Expand All @@ -44,8 +45,8 @@
class DictionaryPageReader implements DictionaryPageReadStore {

private final ParquetFileReader reader;
private final Map<String, ColumnChunkMetaData> columns;
private final Map<String, Optional<DictionaryPage>> dictionaryPageCache;
private final Map<ColumnPath, ColumnChunkMetaData> columns;
private final Map<ColumnPath, Optional<DictionaryPage>> dictionaryPageCache;
private ColumnChunkPageReadStore rowGroup = null;
private ByteBufferReleaser releaser;

Expand All @@ -64,7 +65,7 @@ class DictionaryPageReader implements DictionaryPageReadStore {
releaser = new ByteBufferReleaser(allocator);

for (ColumnChunkMetaData column : block.getColumns()) {
columns.put(column.getPath().toDotString(), column);
columns.put(column.getPath(), column);
}
}

Expand All @@ -87,14 +88,14 @@ public DictionaryPage readDictionaryPage(ColumnDescriptor descriptor) {
return rowGroup.readDictionaryPage(descriptor);
}

String dotPath = String.join(".", descriptor.getPath());
ColumnChunkMetaData column = columns.get(dotPath);
ColumnPath path = ColumnPath.get(descriptor.getPath());
ColumnChunkMetaData column = columns.get(path);
if (column == null) {
throw new ParquetDecodingException("Failed to load dictionary, unknown column: " + dotPath);
throw new ParquetDecodingException("Failed to load dictionary, unknown column: " + path);
}

return dictionaryPageCache
.computeIfAbsent(dotPath, key -> {
.computeIfAbsent(path, key -> {
try {
final DictionaryPage dict = column.hasDictionaryPage() ? reader.readDictionary(column) : null;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,11 +61,14 @@
import org.apache.parquet.bytes.HeapByteBufferAllocator;
import org.apache.parquet.bytes.TrackingByteBufferAllocator;
import org.apache.parquet.column.ColumnDescriptor;
import org.apache.parquet.column.Dictionary;
import org.apache.parquet.column.Encoding;
import org.apache.parquet.column.ParquetProperties;
import org.apache.parquet.column.ParquetProperties.WriterVersion;
import org.apache.parquet.column.page.DataPage;
import org.apache.parquet.column.page.DataPageV2;
import org.apache.parquet.column.page.DictionaryPage;
import org.apache.parquet.column.page.DictionaryPageReadStore;
import org.apache.parquet.column.page.PageReadStore;
import org.apache.parquet.column.page.PageReader;
import org.apache.parquet.column.values.bloomfilter.BloomFilter;
Expand Down Expand Up @@ -335,6 +338,61 @@ public void testParquetFileWithBloomFilter() throws IOException {
}
}

@Test
public void testDictionaryReaderKeepsCollidingDotStringPathsDistinct() 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);
GroupFactory factory = new SimpleGroupFactory(schema);

Path path = newTempPath();
try (ParquetWriter<Group> writer = ExampleParquetWriter.builder(path)
.withAllocator(allocator)
.withConf(conf)
.withDictionaryEncoding(true)
.build()) {
for (int i = 0; i < 100; i++) {
String suffix = (i % 2 == 0) ? "one" : "two";
Group group = factory.newGroup().append("a.b", suffix);
group.addGroup("a").append("b", "nested-" + suffix);
writer.write(group);
}
}

try (ParquetFileReader reader = ParquetFileReader.open(HadoopInputFile.fromPath(path, conf))) {
ColumnDescriptor topLevelDescriptor = schema.getColumnDescription(new String[] {"a.b"});
ColumnDescriptor nestedDescriptor = schema.getColumnDescription(new String[] {"a", "b"});
DictionaryPageReadStore dictionaryReader = reader.getNextDictionaryReader();
DictionaryPage topLevelPage = dictionaryReader.readDictionaryPage(topLevelDescriptor);
DictionaryPage nestedPage = dictionaryReader.readDictionaryPage(nestedDescriptor);
assertThat(topLevelPage).isNotNull();
assertThat(nestedPage).isNotNull();

Dictionary topLevelDictionary = topLevelPage.getEncoding().initDictionary(topLevelDescriptor, topLevelPage);
Set<String> topLevelValues = new HashSet<>();
for (int id = 0; id <= topLevelDictionary.getMaxId(); id++) {
topLevelValues.add(topLevelDictionary.decodeToBinary(id).toStringUsingUTF8());
}
assertThat(topLevelValues).containsExactlyInAnyOrder("one", "two");

Dictionary nestedDictionary = nestedPage.getEncoding().initDictionary(nestedDescriptor, nestedPage);
Set<String> nestedValues = new HashSet<>();
for (int id = 0; id <= nestedDictionary.getMaxId(); id++) {
nestedValues.add(nestedDictionary.decodeToBinary(id).toStringUsingUTF8());
}
assertThat(nestedValues).containsExactlyInAnyOrder("nested-one", "nested-two");
}
}

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