This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 4618d1faaf [common] Bound the file-size memo in both local cache
managers (#9775)
4618d1faaf is described below
commit 4618d1faaf1566c61445bade34dc435e3a0bfb22
Author: YangJie <[email protected]>
AuthorDate: Sun Sep 13 22:50:05 2026 -0400
[common] Bound the file-size memo in both local cache managers (#9775)
---
.../org/apache/paimon/fs/cache/FileSizeMemo.java | 68 ++++++++++++
.../paimon/fs/cache/LocalDiskCacheManager.java | 12 ++-
.../paimon/fs/cache/LocalMemoryCacheManager.java | 19 ++--
.../apache/paimon/fs/cache/CachingFileIOTest.java | 38 +++++++
.../apache/paimon/fs/cache/FileSizeMemoTest.java | 115 +++++++++++++++++++++
5 files changed, 240 insertions(+), 12 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/fs/cache/FileSizeMemo.java
b/paimon-common/src/main/java/org/apache/paimon/fs/cache/FileSizeMemo.java
new file mode 100644
index 0000000000..5691596d36
--- /dev/null
+++ b/paimon-common/src/main/java/org/apache/paimon/fs/cache/FileSizeMemo.java
@@ -0,0 +1,68 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.fs.cache;
+
+import java.util.Iterator;
+import java.util.LinkedHashMap;
+
+/**
+ * Least-recently-used memo of file sizes, bounded by entry count.
+ *
+ * <p>A memo is tiny but there is one per path, so an unbounded map grows with
the number of
+ * distinct files a long-lived process reads. Losing one only costs the extra
{@code getFileStatus}
+ * that would have been made anyway, and the files these caches accept are
immutable, so a re-read
+ * returns the same size.
+ *
+ * <p>Not thread-safe: callers hold their own lock. Access order means {@link
#get} mutates the map,
+ * so even a read has to be inside it.
+ */
+class FileSizeMemo {
+
+ private static final int MAX_ENTRIES = 65536;
+
+ private final LinkedHashMap<String, Long> sizes = new LinkedHashMap<>(64,
0.75f, true);
+
+ /** Read through a method, not the constant: a constant is inlined into
the test's bytecode. */
+ static int maxEntries() {
+ return MAX_ENTRIES;
+ }
+
+ /** Entry count, so a test can observe the bound without reading an entry.
*/
+ int size() {
+ return sizes.size();
+ }
+
+ long get(String filePath) {
+ Long size = sizes.get(filePath);
+ return size != null ? size : -1;
+ }
+
+ void put(String filePath, long size) {
+ sizes.put(filePath, size);
+ Iterator<String> iterator = sizes.keySet().iterator();
+ while (sizes.size() > MAX_ENTRIES && iterator.hasNext()) {
+ iterator.next();
+ iterator.remove();
+ }
+ }
+
+ void invalidate(String filePathPrefix) {
+ sizes.keySet().removeIf(filePath ->
filePath.startsWith(filePathPrefix));
+ }
+}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalDiskCacheManager.java
b/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalDiskCacheManager.java
index d020feac42..4f05b9415f 100644
---
a/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalDiskCacheManager.java
+++
b/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalDiskCacheManager.java
@@ -38,7 +38,6 @@ import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
/** Block-level local disk cache with LRU eviction. Thread-safe. */
public class LocalDiskCacheManager implements LocalCacheManager {
@@ -50,7 +49,7 @@ public class LocalDiskCacheManager implements
LocalCacheManager {
private final long maxSizeBytes;
private final int blockSize;
private final Object lock = new Object();
- private final ConcurrentHashMap<String, Long> fileSizeCache = new
ConcurrentHashMap<>();
+ private final FileSizeMemo fileSizeMemo = new FileSizeMemo();
// LRU-ordered index: key -> size. Access order so get() moves entry to
tail.
private final LinkedHashMap<String, Long> entryIndex;
@@ -240,12 +239,15 @@ public class LocalDiskCacheManager implements
LocalCacheManager {
@Override
public long getFileSize(String filePath) {
- Long size = fileSizeCache.get(filePath);
- return size != null ? size : -1;
+ synchronized (lock) {
+ return fileSizeMemo.get(filePath);
+ }
}
@Override
public void putFileSize(String filePath, long size) {
- fileSizeCache.put(filePath, size);
+ synchronized (lock) {
+ fileSizeMemo.put(filePath, size);
+ }
}
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalMemoryCacheManager.java
b/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalMemoryCacheManager.java
index e92cb88412..541940e668 100644
---
a/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalMemoryCacheManager.java
+++
b/paimon-common/src/main/java/org/apache/paimon/fs/cache/LocalMemoryCacheManager.java
@@ -24,16 +24,18 @@ import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Objects;
-import java.util.concurrent.ConcurrentHashMap;
-/** Block-level in-memory cache with LRU eviction. Thread-safe. */
+/**
+ * In-memory cache with LRU eviction, holding data blocks bounded by total
bytes and a {@link
+ * FileSizeMemo} bounded by entry count. Thread-safe.
+ */
public class LocalMemoryCacheManager implements LocalCacheManager {
private final long maxSizeBytes;
private final int blockSize;
private final Object lock = new Object();
private final LinkedHashMap<BlockKey, byte[]> cache;
- private final ConcurrentHashMap<String, Long> fileSizeCache = new
ConcurrentHashMap<>();
+ private final FileSizeMemo fileSizeMemo = new FileSizeMemo();
private long currentSize;
@@ -80,13 +82,16 @@ public class LocalMemoryCacheManager implements
LocalCacheManager {
@Override
public long getFileSize(String filePath) {
- Long size = fileSizeCache.get(filePath);
- return size != null ? size : -1;
+ synchronized (lock) {
+ return fileSizeMemo.get(filePath);
+ }
}
@Override
public void putFileSize(String filePath, long size) {
- fileSizeCache.put(filePath, size);
+ synchronized (lock) {
+ fileSizeMemo.put(filePath, size);
+ }
}
@Override
@@ -100,8 +105,8 @@ public class LocalMemoryCacheManager implements
LocalCacheManager {
iterator.remove();
}
}
+ fileSizeMemo.invalidate(filePathPrefix);
}
- fileSizeCache.keySet().removeIf(filePath ->
filePath.startsWith(filePathPrefix));
}
private static class BlockKey {
diff --git
a/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
b/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
index 05ddd7cad0..d3ae9d0663 100644
---
a/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
@@ -212,6 +212,44 @@ class CachingFileIOTest {
}
}
+ @Test
+ void fileSizeMemoIsBounded() {
+ // both cache managers keep this memo, and either one is picked purely
by whether
+ // local-cache.dir is set, so the bound has to hold for both
+ assertFileSizeMemoIsBounded(new
LocalMemoryCacheManager(Long.MAX_VALUE, 64));
+ assertFileSizeMemoIsBounded(
+ new LocalDiskCacheManager(
+ tempDir.resolve("memo-bound").toString(),
Long.MAX_VALUE, 64));
+ }
+
+ private static void assertFileSizeMemoIsBounded(LocalCacheManager cache) {
+ // more puts than the bound, so eviction has to run. FileSizeMemoTest
pins the count and
+ // the eviction order; this only checks that the manager routes
through a bounded memo.
+ long entries = FileSizeMemo.maxEntries() + 1024L;
+
+ cache.putFileSize("file-0", 100L);
+ for (long i = 1; i <= entries; i++) {
+ cache.putFileSize("file-" + i, i);
+ }
+
+ assertThat(cache.getFileSize("file-0")).isEqualTo(-1L);
+ assertThat(cache.getFileSize("file-" + entries)).isEqualTo(entries);
+ }
+
+ @Test
+ void memoryCacheInvalidatesFileSizeMemoByPrefix() {
+ // only the memory manager overrides invalidate; the disk one inherits
the no-op default,
+ // which this PR does not change
+ LocalMemoryCacheManager cache = new
LocalMemoryCacheManager(Long.MAX_VALUE, 64);
+ cache.putFileSize("ns/a", 1L);
+ cache.putFileSize("other/a", 2L);
+
+ cache.invalidate("ns/");
+
+ assertThat(cache.getFileSize("ns/a")).isEqualTo(-1L);
+ assertThat(cache.getFileSize("other/a")).isEqualTo(2L);
+ }
+
@Test
void testMetaFileIsCached() throws IOException {
byte[] data = "snapshot data".getBytes();
diff --git
a/paimon-common/src/test/java/org/apache/paimon/fs/cache/FileSizeMemoTest.java
b/paimon-common/src/test/java/org/apache/paimon/fs/cache/FileSizeMemoTest.java
new file mode 100644
index 0000000000..d7f9326494
--- /dev/null
+++
b/paimon-common/src/test/java/org/apache/paimon/fs/cache/FileSizeMemoTest.java
@@ -0,0 +1,115 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.fs.cache;
+
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link FileSizeMemo}. */
+class FileSizeMemoTest {
+
+ @Test
+ void putsAloneBoundTheMemo() {
+ int bound = FileSizeMemo.maxEntries();
+ FileSizeMemo memo = new FileSizeMemo();
+
+ for (int i = 0; i < bound; i++) {
+ memo.put("file-" + i, i);
+ }
+ assertThat(memo.size()).isEqualTo(bound);
+
+ // check the first overflow on its own: a step that evicts the wrong
number of entries
+ // shows up here, where no later put can bring the count back to the
bound
+ memo.put("over-0", 0L);
+ assertThat(memo.size()).isEqualTo(bound);
+
+ for (int i = 1; i < 1024; i++) {
+ memo.put("over-" + i, i);
+ }
+ // no read anywhere above, so the write path is what has to bound it
+ assertThat(memo.size()).isEqualTo(bound);
+ }
+
+ @Test
+ void aReadEntryOutlivesAnUnreadOne() {
+ int bound = FileSizeMemo.maxEntries();
+ // below 4 the four assertions below are not four distinct keys
+ assertThat(bound).isGreaterThanOrEqualTo(4);
+ int read = bound / 2;
+ int unread = bound - read;
+ FileSizeMemo memo = new FileSizeMemo();
+ for (int i = 0; i < bound; i++) {
+ memo.put("file-" + i, i);
+ }
+
+ // reading the older entries makes them the recently used ones
+ for (int i = 0; i < read; i++) {
+ assertThat(memo.get("file-" + i)).isEqualTo(i);
+ }
+ // exactly as many new entries as were left unread, so those are what
eviction takes
+ for (int i = bound; i < bound + unread; i++) {
+ memo.put("file-" + i, i);
+ }
+
+ assertThat(memo.get("file-0")).isEqualTo(0L);
+ assertThat(memo.get("file-" + (read - 1))).isEqualTo(read - 1L);
+ assertThat(memo.get("file-" + read)).isEqualTo(-1L);
+ assertThat(memo.get("file-" + (bound - 1))).isEqualTo(-1L);
+ }
+
+ @Test
+ void puttingAnEntryAgainRefreshesIt() {
+ int bound = FileSizeMemo.maxEntries();
+ // at a bound of 1 the loop below never runs, so nothing would pin the
position half
+ assertThat(bound).isGreaterThanOrEqualTo(2);
+ FileSizeMemo memo = new FileSizeMemo();
+ for (int i = 0; i < bound; i++) {
+ memo.put("file-" + i, i);
+ }
+
+ memo.put("file-0", 100L);
+ for (int i = bound; i < bound + bound - 1; i++) {
+ memo.put("file-" + i, i);
+ }
+
+ // the re-put carried both the newer value and the newer position
+ assertThat(memo.get("file-0")).isEqualTo(100L);
+ assertThat(memo.get("file-1")).isEqualTo(-1L);
+ }
+
+ @Test
+ void invalidateRemovesOnlyTheMatchingPrefix() {
+ // the fixture below holds four entries, and none of them may be
evicted
+ assertThat(FileSizeMemo.maxEntries()).isGreaterThanOrEqualTo(4);
+ FileSizeMemo memo = new FileSizeMemo();
+ memo.put("/a/one", 1L);
+ memo.put("/a/two", 2L);
+ memo.put("/b/three", 3L);
+ // carries the prefix, but not at the front
+ memo.put("/b/a/four", 4L);
+
+ memo.invalidate("/a/");
+
+ assertThat(memo.get("/a/one")).isEqualTo(-1L);
+ assertThat(memo.get("/a/two")).isEqualTo(-1L);
+ assertThat(memo.get("/b/three")).isEqualTo(3L);
+ assertThat(memo.get("/b/a/four")).isEqualTo(4L);
+ }
+}