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 @@ -36,6 +36,7 @@
import java.util.Map.Entry;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
import java.util.zip.CRC32;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataInputStream;
Expand Down Expand Up @@ -1702,29 +1703,31 @@ public void appendRowGroup(SeekableInputStream from, BlockMetaData rowGroup, boo
withAbortOnFailure(() -> {
startBlock(rowGroup.getRowCount());

Map<String, ColumnChunkMetaData> columnsToCopy = new HashMap<String, ColumnChunkMetaData>();
Map<ColumnPath, ColumnChunkMetaData> columnsToCopy = new HashMap<ColumnPath, ColumnChunkMetaData>();
for (ColumnChunkMetaData chunk : rowGroup.getColumns()) {
columnsToCopy.put(chunk.getPath().toDotString(), chunk);
columnsToCopy.put(chunk.getPath(), chunk);
}

List<ColumnChunkMetaData> columnsInOrder = new ArrayList<ColumnChunkMetaData>();

for (ColumnDescriptor descriptor : schema.getColumns()) {
String path = ColumnPath.get(descriptor.getPath()).toDotString();
ColumnPath path = ColumnPath.get(descriptor.getPath());
ColumnChunkMetaData chunk = columnsToCopy.remove(path);
if (chunk != null) {
columnsInOrder.add(chunk);
} else {
throw new IllegalArgumentException(
String.format("Missing column '%s', cannot copy row group: %s", path, rowGroup));
throw new IllegalArgumentException(String.format(
"Missing column '%s', cannot copy row group: %s", path.toDotString(), rowGroup));
}
}

// complain if some columns would be dropped and that's not okay
if (!dropColumns && !columnsToCopy.isEmpty()) {
throw new IllegalArgumentException(String.format(
"Columns cannot be copied (missing from target schema): %s",
String.join(", ", columnsToCopy.keySet())));
columnsToCopy.keySet().stream()
.map(ColumnPath::toDotString)
.collect(Collectors.joining(", "))));
}

// copy the data for all chunks
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,43 @@ public void testBasicBehavior() throws IOException {
assertThat(expected).as("All records should be present").isEmpty();
}

@Test
public void testAppendKeepsCollidingDotStringPathsDistinct() throws IOException {
MessageType schema = Types.buildMessage()
.required(BINARY)
.as(UTF8)
.named("a.b")
.requiredGroup()
.required(BINARY)
.as(UTF8)
.named("b")
.named("a")
.named("AppendCollisionTest");
SimpleGroupFactory factory = new SimpleGroupFactory(schema);
Group expected = factory.newGroup().append("a.b", "top-level");
expected.addGroup("a").append("b", "nested");

Path source = newTemp();
try (ParquetWriter<Group> sourceWriter =
ExampleParquetWriter.builder(source).withType(schema).build()) {
sourceWriter.write(expected);
}

Path copied = newTemp();
ParquetFileWriter copyWriter = new ParquetFileWriter(CONF, schema, copied);
copyWriter.start();
copyWriter.appendFile(CONF, source);
copyWriter.end(EMPTY_METADATA);

try (ParquetReader<Group> reader =
ParquetReader.builder(new GroupReadSupport(), copied).build()) {
Group actual = reader.read();
assertThat(actual.getString("a.b", 0)).isEqualTo("top-level");
assertThat(actual.getGroup("a", 0).getString("b", 0)).isEqualTo("nested");
assertThat(reader.read()).isNull();
}
}

/**
* This test is similar to {@link #testBasicBehavior()} only that it uses static files generated by a previous release
* (1.11.1). This test is to validate the fix of PARQUET-2027.
Expand Down