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

Reply via email to