This is an automated email from the ASF dual-hosted git repository.

danny0405 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 5550b4585631 perf(flink): streamline cache eviction and metadata 
refresh (#20134)
5550b4585631 is described below

commit 5550b458563131bd0bfb331ee364944394fa5ddd
Author: Danny Chan <[email protected]>
AuthorDate: Tue Sep 29 17:26:42 2026 +0800

    perf(flink): streamline cache eviction and metadata refresh (#20134)
    
    * perf(flink): streamline record-level index cache cleanup
    * perf(flink): reload coordinator metadata in place
---
 .../hudi/sink/StreamWriteOperatorCoordinator.java  |  4 +--
 .../index/GlobalRecordLevelIndexBackend.java       |  1 -
 .../partitioner/index/RecordLevelIndexBackend.java | 33 ++++++++------------
 .../sink/TestStreamWriteOperatorCoordinator.java   | 25 ---------------
 .../index/TestRecordLevelIndexBackend.java         | 36 ++++++++++++++++++++++
 5 files changed, 50 insertions(+), 49 deletions(-)

diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
index 3100cb5f3cac..1cdc59a8286a 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
@@ -518,8 +518,8 @@ public class StreamWriteOperatorCoordinator
   }
 
   private String startInstant() {
-    // Refresh table properties and index definitions as well as the timeline 
before validating the new write.
-    this.metaClient = HoodieTableMetaClient.reload(this.metaClient);
+    // Refresh table properties and the timeline before validating the new 
write.
+    this.metaClient.reload();
     // Validate the write and refresh the last txn metadata.
     this.writeClient.preTxn(tableState.operationType, this.metaClient);
     // put the assignment in front of metadata generation,
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java
index 791c064007a7..1554e99c500d 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java
@@ -169,7 +169,6 @@ public class GlobalRecordLevelIndexBackend implements 
MinibatchIndexBackend {
     // the latest completed checkpoint id is used as the minimum checkpoint id,
     // since the streaming write operator always uses previous checkpoint id 
to request the new instant.
     
recordIndexCache.markAsEvictable(inflightInstants.keySet().stream().min(Long::compareTo).orElse(completedCheckpointID));
-    this.metaClient.reloadActiveTimeline();
     reloadMetadataTable();
   }
 
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/RecordLevelIndexBackend.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/RecordLevelIndexBackend.java
index 3d19c2061ceb..22b8baddde4f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/RecordLevelIndexBackend.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/RecordLevelIndexBackend.java
@@ -154,7 +154,6 @@ public class RecordLevelIndexBackend implements 
PartitionedIndexBackend {
   public void onCheckpointComplete(Correspondent correspondent, long 
completedCheckpointId) {
     Map<Long, String> inflightInstants = 
correspondent.requestInflightInstants();
     
updateEvictableCkp(inflightInstants.keySet().stream().min(Long::compareTo).orElse(completedCheckpointId));
-    metaClient.reloadActiveTimeline();
     reloadMetadataTable();
   }
 
@@ -291,27 +290,19 @@ public class RecordLevelIndexBackend implements 
PartitionedIndexBackend {
 
   @VisibleForTesting
   void cleanIfNecessary(long nextCacheSize, String protectedPartitionPath) {
-    while (getCurrentHeapSize() + nextCacheSize > maxCacheSizeInBytes) {
-      boolean cleaned = false;
-      Iterator<Map.Entry<String, BucketCache>> iterator = 
partitionBucketCaches.entrySet().iterator();
-      while (iterator.hasNext()) {
-        Map.Entry<String, BucketCache> entry = iterator.next();
-        if (entry.getKey().equals(protectedPartitionPath)) {
-          continue;
-        }
-        BucketCache cache = entry.getValue();
-        if (cache.lastUpdatedCheckpoint < minRetainedCheckpointId) {
-          cache.close();
-          iterator.remove();
-          cleaned = true;
-          log.info("Evict partitioned RLI cache for partition {}", 
entry.getKey());
-          break;
-        }
+    long currentHeapSize = getCurrentHeapSize();
+    Iterator<Map.Entry<String, BucketCache>> iterator = 
partitionBucketCaches.entrySet().iterator();
+    while (currentHeapSize + nextCacheSize > maxCacheSizeInBytes && 
iterator.hasNext()) {
+      Map.Entry<String, BucketCache> entry = iterator.next();
+      if (entry.getKey().equals(protectedPartitionPath)) {
+        continue;
       }
-      if (!cleaned) {
-        // All remaining partition caches are either protected or too recent 
to evict safely.
-        // Returning avoids retrying the same scan without making progress.
-        return;
+      BucketCache cache = entry.getValue();
+      if (cache.lastUpdatedCheckpoint < minRetainedCheckpointId) {
+        currentHeapSize -= cache.getHeapSize();
+        cache.close();
+        iterator.remove();
+        log.info("Evict partitioned RLI cache for partition {}", 
entry.getKey());
       }
     }
   }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
index 12adab2cc3e0..7dd73af9dfaa 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
@@ -25,7 +25,6 @@ import org.apache.hudi.client.heartbeat.HoodieHeartbeatClient;
 import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.model.HoodieCommitMetadata;
 import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
-import org.apache.hudi.common.model.HoodieIndexDefinition;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.model.HoodieWriteStat;
 import org.apache.hudi.common.model.MetaFieldsMode;
@@ -47,9 +46,7 @@ import org.apache.hudi.configuration.FlinkOptions;
 import org.apache.hudi.configuration.HadoopConfigurations;
 import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.exception.MissingSchemaFieldException;
-import org.apache.hudi.exception.SchemaCompatibilityException;
 import org.apache.hudi.hadoop.fs.HadoopFSUtils;
-import org.apache.hudi.metadata.HoodieIndexVersion;
 import org.apache.hudi.metadata.HoodieTableMetadata;
 import org.apache.hudi.metadata.MetadataPartitionType;
 import org.apache.hudi.sink.event.Correspondent;
@@ -218,28 +215,6 @@ public class TestStreamWriteOperatorCoordinator {
     assertInstantCreationFails(conf, HoodieException.class, 
HoodieTableConfig.META_FIELDS_MODE.key());
   }
 
-  @Test
-  void testNewSecondaryIndexIsValidatedBeforeInstantIsPublished() throws 
Exception {
-    Configuration conf = 
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
-    HoodieWriteConfig writeConfig = coordinator.getWriteClient().getConfig();
-    HoodieTableMetaClient metaClient = StreamerUtil.createMetaClient(conf);
-    HoodieCommitMetadata metadata = new HoodieCommitMetadata();
-    metadata.addMetadata(HoodieCommitMetadata.SCHEMA_KEY, 
writeConfig.getSchema());
-    HoodieTestTable.of(metaClient).addCommit("001", Option.of(metadata));
-    // Simulate an index created after the coordinator has loaded its meta 
client.
-    metaClient.buildIndexDefinition(HoodieIndexDefinition.newBuilder()
-        .withIndexName("secondary_index_age")
-        .withIndexType("secondary_index")
-        .withSourceFields(Collections.singletonList("age"))
-        .withVersion(HoodieIndexVersion.V1)
-        .build());
-    String evolvedSchema = writeConfig.getSchema().replace("\"int\"", 
"\"long\"");
-    assertNotEquals(writeConfig.getSchema(), evolvedSchema);
-    writeConfig.setSchema(evolvedSchema);
-
-    assertInstantCreationFails(conf, SchemaCompatibilityException.class, 
"secondary_index_age");
-  }
-
   private void assertInstantCreationFails(Configuration conf, Class<? extends 
Throwable> causeType, String message) throws Exception {
     Correspondent.InstantTimeResponse response = 
CoordinationResponseSerDe.unwrap(
         
coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(1))
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestRecordLevelIndexBackend.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestRecordLevelIndexBackend.java
index 83ee81f80ecc..ad7b51bcd70b 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestRecordLevelIndexBackend.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestRecordLevelIndexBackend.java
@@ -33,6 +33,8 @@ import org.junit.jupiter.api.io.TempDir;
 import org.mockito.Mockito;
 
 import java.io.File;
+import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
@@ -43,7 +45,10 @@ import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.atMost;
 import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 /**
@@ -93,6 +98,37 @@ public class TestRecordLevelIndexBackend {
     }
   }
 
+  @Test
+  public void testLazyEvictMultiplePartitionsWithoutRepeatedHeapScans() throws 
Exception {
+    try (RecordLevelIndexBackend backend = createBackend()) {
+      Map<String, ExternalSpillableMap<String, Integer>> maps = new 
HashMap<>();
+      for (String partition : Arrays.asList("protected", "old1", "recent", 
"old2", "old3", "retained")) {
+        ExternalSpillableMap<String, Integer> map = mapWithStorage(ONE_MB / 4);
+        maps.put(partition, map);
+        backend.getPartitionBucketCaches().put(partition,
+            backend.newBucketCache(map, partition.equals("recent") ? 2L : 1L));
+        doAnswer(invocation -> {
+          when(map.getCurrentInMemoryMapSize()).thenReturn(0L);
+          return null;
+        }).when(map).close();
+      }
+      backend.onCheckpointComplete(new 
TestCorrespondent(Collections.singletonMap(2L, "002")), 3L);
+
+      backend.cleanIfNecessary(ONE_MB / 4, "protected");
+
+      assertEquals(Arrays.asList("protected", "recent", "retained"),
+          new ArrayList<>(backend.getPartitionBucketCaches().keySet()));
+      for (Map.Entry<String, ExternalSpillableMap<String, Integer>> entry : 
maps.entrySet()) {
+        verify(entry.getValue(), atMost(2)).getCurrentInMemoryMapSize();
+        if (backend.getPartitionBucketCaches().containsKey(entry.getKey())) {
+          verify(entry.getValue(), never()).close();
+        } else {
+          verify(entry.getValue()).close();
+        }
+      }
+    }
+  }
+
   @Test
   public void testLazyEvictKeepsRecentPartition() throws Exception {
     try (RecordLevelIndexBackend backend = createBackend()) {

Reply via email to