This is an automated email from the ASF dual-hosted git repository.
voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new d7360762e1fc fix(core): keep plan generator extra metadata when
scheduling compaction (#19941)
d7360762e1fc is described below
commit d7360762e1fc2320fa7afe1101068a7e371067fa
Author: Y Ethan Guo <[email protected]>
AuthorDate: Mon Sep 14 23:02:46 2026 -0700
fix(core): keep plan generator extra metadata when scheduling compaction
(#19941)
* fix(core): keep plan generator extra metadata when scheduling compaction
The table service client enriches the extra metadata passed to compaction
scheduling, so the option is never empty, and the scheduler replaced the
plan's extra metadata with it. Entries the plan generator recorded were
dropped on every schedule. Merge the two instead, letting the generator's
entries take precedence.
* Describe the merged plan metadata generically
---
.../compact/ScheduleCompactionActionExecutor.java | 21 ++++++++-
.../TestScheduleCompactionActionExecutor.java | 51 ++++++++++++++++++++++
.../table/action/compact/TestHoodieCompactor.java | 48 ++++++++++++++++++++
3 files changed, 119 insertions(+), 1 deletion(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/ScheduleCompactionActionExecutor.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/ScheduleCompactionActionExecutor.java
index 6f03b69dcb04..5b8a11cdcd54 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/ScheduleCompactionActionExecutor.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/ScheduleCompactionActionExecutor.java
@@ -33,6 +33,7 @@ import org.apache.hudi.common.util.ConfigUtils;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.ReflectionUtils;
import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.common.util.VisibleForTesting;
import org.apache.hudi.common.util.collection.Pair;
import org.apache.hudi.config.HoodieCompactionConfig;
import org.apache.hudi.config.HoodieWriteConfig;
@@ -47,6 +48,7 @@ import lombok.extern.slf4j.Slf4j;
import javax.annotation.Nullable;
import java.io.IOException;
+import java.util.HashMap;
import java.util.Map;
import static
org.apache.hudi.common.table.timeline.InstantComparison.GREATER_THAN;
@@ -106,7 +108,7 @@ public class ScheduleCompactionActionExecutor<T, I, K, O>
extends BaseTableServi
HoodieCompactionPlan plan = scheduleCompaction();
Option<HoodieCompactionPlan> option = Option.empty();
if (plan != null && nonEmpty(plan.getOperations())) {
- extraMetadata.ifPresent(plan::setExtraMetadata);
+ mergeExtraMetadata(plan, extraMetadata);
if (operationType.equals(WriteOperationType.COMPACT)) {
HoodieInstant compactionInstant =
instantGenerator.createNewInstant(HoodieInstant.State.REQUESTED,
HoodieTimeline.COMPACTION_ACTION, instantTime);
@@ -122,6 +124,23 @@ public class ScheduleCompactionActionExecutor<T, I, K, O>
extends BaseTableServi
return option;
}
+ /**
+ * Adds the caller-provided extra metadata to the plan without discarding
what the plan
+ * generator already recorded there. Those entries carry state only the
generator can
+ * reconstruct, so they take precedence on a key collision; every other
caller entry is kept.
+ */
+ @VisibleForTesting
+ static void mergeExtraMetadata(HoodieCompactionPlan plan, Option<Map<String,
String>> extraMetadata) {
+ if (!extraMetadata.isPresent()) {
+ return;
+ }
+ Map<String, String> merged = new HashMap<>(extraMetadata.get());
+ if (plan.getExtraMetadata() != null) {
+ merged.putAll(plan.getExtraMetadata());
+ }
+ plan.setExtraMetadata(merged);
+ }
+
@Nullable
private HoodieCompactionPlan scheduleCompaction() {
log.info("Checking if compaction needs to be run on {}",
config.getBasePath());
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
index 4d9ba937a4eb..f7ad94322fe3 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
@@ -18,6 +18,7 @@
package org.apache.hudi.table.action.compact;
+import org.apache.hudi.avro.model.HoodieCompactionPlan;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.engine.HoodieLocalEngineContext;
@@ -39,6 +40,9 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import java.io.IOException;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
import static org.apache.hudi.hadoop.fs.HadoopFSUtils.getStorageConf;
import static org.mockito.Mockito.when;
@@ -65,6 +69,53 @@ class TestScheduleCompactionActionExecutor extends
HoodieCommonTestHarness {
Assertions.assertEquals(1, TestCompactionPlanGenerator.getCount());
}
+ @Test
+ void testMergeExtraMetadataKeepsPlanGeneratorEntries() {
+ HoodieCompactionPlan plan = planWithExtraMetadata(generatorMetadata());
+ Map<String, String> provided = new HashMap<>();
+ provided.put("caller.key", "caller.value");
+ ScheduleCompactionActionExecutor.mergeExtraMetadata(plan,
Option.of(provided));
+ Map<String, String> expected = new HashMap<>(generatorMetadata());
+ expected.putAll(provided);
+ Assertions.assertEquals(expected, plan.getExtraMetadata());
+ }
+
+ @Test
+ void testMergeExtraMetadataGeneratorWinsOnCollision() {
+ HoodieCompactionPlan plan = planWithExtraMetadata(generatorMetadata());
+ ScheduleCompactionActionExecutor.mergeExtraMetadata(plan,
Option.of(Collections.singletonMap("generator.mode", "other")));
+ Assertions.assertEquals(generatorMetadata(), plan.getExtraMetadata());
+ }
+
+ @Test
+ void testMergeExtraMetadataWithoutCallerEntriesLeavesPlanUntouched() {
+ HoodieCompactionPlan plan = planWithExtraMetadata(generatorMetadata());
+ ScheduleCompactionActionExecutor.mergeExtraMetadata(plan, Option.empty());
+ Assertions.assertEquals(generatorMetadata(), plan.getExtraMetadata());
+ }
+
+ @Test
+ void testMergeExtraMetadataOntoPlanWithoutEntries() {
+ HoodieCompactionPlan plan = planWithExtraMetadata(null);
+ Map<String, String> provided = Collections.singletonMap("caller.key",
"caller.value");
+ ScheduleCompactionActionExecutor.mergeExtraMetadata(plan,
Option.of(provided));
+ Assertions.assertEquals(provided, plan.getExtraMetadata());
+ }
+
+ private static Map<String, String> generatorMetadata() {
+ Map<String, String> metadata = new HashMap<>();
+ metadata.put("generator.mode", "planned");
+ metadata.put("generator.state", "[\"a\",\"b\"]");
+ return metadata;
+ }
+
+ private static HoodieCompactionPlan planWithExtraMetadata(Map<String,
String> extraMetadata) {
+ return HoodieCompactionPlan.newBuilder()
+ .setOperations(Collections.emptyList())
+ .setExtraMetadata(extraMetadata)
+ .build();
+ }
+
public static class TestCompactionPlanGenerator<T extends
HoodieRecordPayload, I, K, O> extends HoodieCompactionPlanGenerator<T, I, K, O>
{
private static int count = 0;
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/action/compact/TestHoodieCompactor.java
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/action/compact/TestHoodieCompactor.java
index 7cdecd119492..31a63f0dd976 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/action/compact/TestHoodieCompactor.java
+++
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/table/action/compact/TestHoodieCompactor.java
@@ -26,9 +26,11 @@ import org.apache.hudi.client.WriteClientTestUtils;
import org.apache.hudi.common.config.HoodieMemoryConfig;
import org.apache.hudi.common.config.HoodieStorageConfig;
import org.apache.hudi.common.config.metrics.HoodieMetricsConfig;
+import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.model.FileSlice;
import org.apache.hudi.common.model.HoodieCommitMetadata;
import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordPayload;
import org.apache.hudi.common.model.HoodieTableType;
import org.apache.hudi.common.model.HoodieWriteStat;
import org.apache.hudi.common.model.TableServiceType;
@@ -52,7 +54,9 @@ import
org.apache.hudi.index.bloom.SparkHoodieBloomIndexHelper;
import org.apache.hudi.metrics.HoodieMetrics;
import org.apache.hudi.table.HoodieSparkTable;
import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.action.BaseTableServicePlanActionExecutor;
import org.apache.hudi.table.action.HoodieWriteMetadata;
+import
org.apache.hudi.table.action.compact.plan.generators.HoodieCompactionPlanGenerator;
import
org.apache.hudi.table.action.compact.strategy.PartitionRegexBasedCompactionStrategy;
import
org.apache.hudi.table.action.compact.strategy.SmallBoundedIOCompactionStrategy;
import org.apache.hudi.testutils.HoodieSparkClientTestHarness;
@@ -70,6 +74,7 @@ import java.io.IOException;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import java.util.Map;
import java.util.SortedMap;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -258,6 +263,49 @@ public class TestHoodieCompactor extends
HoodieSparkClientTestHarness {
}
}
+ @Test
+ public void testScheduleCompactionKeepsPlanGeneratorExtraMetadata() throws
Exception {
+ HoodieWriteConfig config = getConfigBuilder()
+
.withCompactionConfig(HoodieCompactionConfig.newBuilder().withMaxNumDeltaCommitsBeforeCompaction(1).build())
+ .withProps(Collections.singletonMap(
+ HoodieCompactionConfig.COMPACTION_PLAN_GENERATOR.key(),
ExtraMetadataCompactionPlanGenerator.class.getName()))
+ .build();
+ try (SparkRDDWriteClient writeClient = getHoodieWriteClient(config)) {
+ String newCommitTime = "100";
+ WriteClientTestUtils.startCommitWithTime(writeClient, newCommitTime);
+ List<HoodieRecord> records = dataGen.generateInserts(newCommitTime, 100);
+ writeClient.commit(newCommitTime,
writeClient.insert(jsc.parallelize(records, 1), newCommitTime));
+ updateRecords(config, "101", records);
+
+ Option<String> compactionInstant =
writeClient.scheduleCompaction(Option.of(Collections.singletonMap("caller.marker",
"true")));
+ assertTrue(compactionInstant.isPresent());
+ metaClient.reloadActiveTimeline();
+ Map<String, String> extraMetadata =
CompactionUtils.getCompactionPlan(metaClient,
compactionInstant.get()).getExtraMetadata();
+ // The generator's entry and the caller's entry must both be on the
persisted plan.
+ assertEquals("true",
extraMetadata.get(ExtraMetadataCompactionPlanGenerator.MARKER_KEY));
+ assertEquals("true", extraMetadata.get("caller.marker"));
+ }
+ }
+
+ /**
+ * Plan generator that records its own state in the plan's extra metadata.
+ */
+ public static class ExtraMetadataCompactionPlanGenerator<T extends
HoodieRecordPayload, I, K, O>
+ extends HoodieCompactionPlanGenerator<T, I, K, O> {
+ static final String MARKER_KEY = "generator.marker";
+
+ public ExtraMetadataCompactionPlanGenerator(HoodieTable table,
HoodieEngineContext engineContext,
+ HoodieWriteConfig writeConfig,
BaseTableServicePlanActionExecutor executor) {
+ super(table, engineContext, writeConfig, executor);
+ }
+
+ @Override
+ protected Map<String, String>
getExtraMetadata(List<HoodieCompactionOperation>
operationsBeforeApplyingStrategy,
+ HoodieCompactionPlan
compactionPlan) {
+ return Collections.singletonMap(MARKER_KEY, "true");
+ }
+ }
+
@Test
public void testSpillingWhenCompaction() throws Exception {
// insert 100 records