Skip to content
Merged
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 @@ -94,7 +94,14 @@ public static int repair(
private static List<Map<String, String>> listFilesystemPartitionSpecs(FormatTable formatTable) {
// Raw directory names rather than the table scan: the scan casts each value to its column
// type and back (month=01 -> 1), producing specs that no longer name the real directory,
// while the write path registers the raw value.
// while the write path registers the raw value. That is also why no partition type is
// passed below: it would re-enable the cast this discovery deliberately avoids.
//
// The default partition name is still needed. In a value-only layout the null partition is
// a bare "__DEFAULT_PARTITION__" directory, which the generic hidden-directory rule ("_"
// prefix) skips unless the listing knows the name is meaningful. Without it a repair never
// registers the null partition, and a SYNC/DROP sees it as registered-but-deleted and
// unregisters a partition that still holds data.
boolean onlyValueInPath =
new CoreOptions(formatTable.options()).formatTablePartitionOnlyValueInPath();
List<Pair<LinkedHashMap<String, String>, Path>> found =
Expand All @@ -103,7 +110,10 @@ private static List<Map<String, String>> listFilesystemPartitionSpecs(FormatTabl
new Path(formatTable.location()),
formatTable.partitionKeys().size(),
formatTable.partitionKeys(),
onlyValueInPath);
onlyValueInPath,
null,
null,
formatTable.defaultPartName());
List<Map<String, String>> specs = new ArrayList<>(found.size());
for (Pair<LinkedHashMap<String, String>, Path> pair : found) {
PartitionPathUtils.validatePartitionSpecForPath(pair.getKey(), onlyValueInPath);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,111 @@ void repairRegistersRawDirectoryValuesWithoutCastingThem() throws Exception {
assertThat(catalog.createIgnoreFlags).containsExactly(true);
}

@Test
void repairAddsTheNullPartitionDirectoryInValueOnlyLayout() throws Exception {
// A value-only layout writes the null partition as a bare __DEFAULT_PARTITION__ directory,
// which the generic hidden-directory rule ("_" prefix) would swallow.
writeDataFile(tempDir.resolve("20260701"));
writeDataFile(tempDir.resolve("__DEFAULT_PARTITION__"));

RecordingPartitionManager catalog = new RecordingPartitionManager();
PaimonFormatTable sparkTable =
new PaimonFormatTable(formatTable(tempDir.toUri().toString(), true, catalog));

int applied = FormatTablePartitionRepair.repair(sparkTable, true, false);

assertThat(applied).isEqualTo(2);
assertThat(catalog.createdPartitions)
.containsExactly(
Arrays.asList(spec("dt", "20260701"), spec("dt", "__DEFAULT_PARTITION__")));
assertThat(catalog.droppedPartitions).isEmpty();
}

@Test
void repairKeepsTheRegisteredNullPartitionInValueOnlyLayout() throws Exception {
// Both directories exist, so a SYNC must be a no-op. Missing the null partition on the
// filesystem side makes it look "registered but deleted" and silently unregisters live
// data.
writeDataFile(tempDir.resolve("20260701"));
writeDataFile(tempDir.resolve("__DEFAULT_PARTITION__"));

RecordingPartitionManager catalog = new RecordingPartitionManager();
catalog.register(
Arrays.asList(spec("dt", "20260701"), spec("dt", "__DEFAULT_PARTITION__")));
PaimonFormatTable sparkTable =
new PaimonFormatTable(formatTable(tempDir.toUri().toString(), true, catalog));

int applied = FormatTablePartitionRepair.repair(sparkTable, true, true);

assertThat(applied).isZero();
assertThat(catalog.droppedPartitions).isEmpty();
assertThat(catalog.createdPartitions).isEmpty();
}

@Test
void repairReadsTheDefaultPartitionNameFromTableOptions() throws Exception {
// Pins that the rescued name comes from partition.default-name rather than a literal:
// hardcoding "__DEFAULT_PARTITION__" reproduces the bug for anyone who overrides it.
writeDataFile(tempDir.resolve("20260701"));
writeDataFile(tempDir.resolve("__MY_NULL__"));

RecordingPartitionManager catalog = new RecordingPartitionManager();
Map<String, String> extra = new LinkedHashMap<>();
extra.put(CoreOptions.PARTITION_DEFAULT_NAME.key(), "__MY_NULL__");
PaimonFormatTable sparkTable =
new PaimonFormatTable(
formatTable(tempDir.toUri().toString(), true, catalog, extra));

int applied = FormatTablePartitionRepair.repair(sparkTable, true, false);

assertThat(applied).isEqualTo(2);
assertThat(catalog.createdPartitions)
.containsExactly(Arrays.asList(spec("dt", "20260701"), spec("dt", "__MY_NULL__")));
}

@Test
void repairDescendsIntoANullPartitionSubtreeInValueOnlyLayout() throws Exception {
// A null value on a non-leaf level hides the whole subtree, not just one directory:
// listStatusRecursively applies the same hidden-name rule while descending.
writeDataFile(tempDir.resolve("20260701").resolve("01"));
writeDataFile(tempDir.resolve("__DEFAULT_PARTITION__").resolve("01"));

RecordingPartitionManager catalog = new RecordingPartitionManager();
PaimonFormatTable sparkTable =
new PaimonFormatTable(twoLevelValueOnlyTable(tempDir.toUri().toString(), catalog));

int applied = FormatTablePartitionRepair.repair(sparkTable, true, false);

Map<String, String> real = new LinkedHashMap<>();
real.put("dt", "20260701");
real.put("month", "01");
Map<String, String> nullDt = new LinkedHashMap<>();
nullDt.put("dt", "__DEFAULT_PARTITION__");
nullDt.put("month", "01");
assertThat(applied).isEqualTo(2);
assertThat(catalog.createdPartitions).containsExactly(Arrays.asList(real, nullDt));
}

@Test
void repairKeepsAnUnderscoreValueInTheKeyValueLayout() throws Exception {
// The hidden-name rule reads the directory name, and in a key=value layout that name is
// "dt=_abc" - the underscore sits on the value, not on the first character. Pins that the
// value-only defect does not extend to the default layout.
writeDataFile(tempDir.resolve("dt=20260701"));
writeDataFile(tempDir.resolve("dt=_abc"));

RecordingPartitionManager catalog = new RecordingPartitionManager();
catalog.register(Arrays.asList(spec("dt", "20260701"), spec("dt", "_abc")));
PaimonFormatTable sparkTable =
new PaimonFormatTable(formatTable(tempDir.toUri().toString(), catalog));

int applied = FormatTablePartitionRepair.repair(sparkTable, true, true);

assertThat(applied).isZero();
assertThat(catalog.droppedPartitions).isEmpty();
assertThat(catalog.createdPartitions).isEmpty();
}

@Test
void repairRejectsUnsafeValueOnlyDirectoryBeforeCatalogMutation() throws Exception {
Files.createDirectories(tempDir.resolve("%2E%2E"));
Expand Down Expand Up @@ -501,6 +606,14 @@ void repairWithoutMeasuringKeepsTheSpecOnlyRegistration() throws Exception {
assertThat(catalog.replaceFlags).containsExactly(false);
}

private static void writeDataFile(java.nio.file.Path partitionDirectory) throws IOException {
Files.createDirectories(partitionDirectory);
Files.write(
partitionDirectory.resolve("data.csv"),
Collections.singletonList("1"),
StandardCharsets.UTF_8);
}

private static Map<String, String> spec(String key, String value) {
Map<String, String> spec = new LinkedHashMap<>();
spec.put(key, value);
Expand All @@ -511,6 +624,45 @@ private static FormatTable formatTable(String location, FormatTablePartitionMana
return formatTable(location, false, catalog);
}

private static FormatTable formatTable(
String location,
boolean onlyValueInPath,
FormatTablePartitionManager catalog,
Map<String, String> extraOptions) {
RowType rowType =
RowType.builder()
.field("id", DataTypes.INT())
.field("dt", DataTypes.STRING())
.build();
return build(
LocalFileIO.create(),
location,
rowType,
Collections.singletonList("dt"),
onlyValueInPath,
catalog,
extraOptions);
}

/** Two STRING partition keys in a value-only layout, so a null value can sit on a non-leaf. */
private static FormatTable twoLevelValueOnlyTable(
String location, FormatTablePartitionManager catalog) {
RowType rowType =
RowType.builder()
.field("id", DataTypes.INT())
.field("dt", DataTypes.STRING())
.field("month", DataTypes.STRING())
.build();
return build(
LocalFileIO.create(),
location,
rowType,
Arrays.asList("dt", "month"),
true,
catalog,
Collections.emptyMap());
}

private static FormatTable formatTable(
String location, boolean onlyValueInPath, FormatTablePartitionManager catalog) {
return formatTable(LocalFileIO.create(), location, onlyValueInPath, catalog);
Expand Down Expand Up @@ -554,11 +706,30 @@ private static FormatTable build(
List<String> partitionKeys,
boolean onlyValueInPath,
FormatTablePartitionManager catalog) {
return build(
fileIO,
location,
rowType,
partitionKeys,
onlyValueInPath,
catalog,
Collections.emptyMap());
}

private static FormatTable build(
FileIO fileIO,
String location,
RowType rowType,
List<String> partitionKeys,
boolean onlyValueInPath,
FormatTablePartitionManager catalog,
Map<String, String> extraOptions) {
Map<String, String> options = new LinkedHashMap<>();
options.put(CoreOptions.METASTORE_PARTITIONED_TABLE.key(), "true");
options.put(
CoreOptions.FORMAT_TABLE_PARTITION_ONLY_VALUE_IN_PATH.key(),
Boolean.toString(onlyValueInPath));
options.putAll(extraOptions);
return FormatTable.builder()
.fileIO(fileIO)
.identifier(Identifier.create("db", "t"))
Expand Down
Loading