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()) {