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

jojochuang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/master by this push:
     new 61cb7834055 HDDS-16229. Add RocksDB multiGetSkipCache for batched 
table reads. (#11068)
61cb7834055 is described below

commit 61cb7834055ddb118631c2eeefb1d0f242e4bfc6
Author: SaketaChalamchala <[email protected]>
AuthorDate: Thu Sep 10 08:45:10 2026 -0700

    HDDS-16229. Add RocksDB multiGetSkipCache for batched table reads. (#11068)
---
 .../ozone/container/metadata/DatanodeTable.java    |   5 +
 .../org/apache/hadoop/hdds/utils/db/RDBTable.java  |  24 ++++
 .../apache/hadoop/hdds/utils/db/RocksDatabase.java |  42 ++++++
 .../org/apache/hadoop/hdds/utils/db/Table.java     |  22 +++
 .../apache/hadoop/hdds/utils/db/TypedTable.java    | 148 +++++++++++++++++++++
 .../hadoop/hdds/utils/db/InMemoryTestTable.java    |  13 ++
 .../apache/hadoop/hdds/utils/db/TestRDBTable.java  |  23 ++++
 .../hadoop/hdds/utils/db/TestRDBTableStore.java    |  99 ++++++++++++++
 pom.xml                                            |   1 +
 9 files changed, 377 insertions(+)

diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/metadata/DatanodeTable.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/metadata/DatanodeTable.java
index fd73f6547c7..23c3afe5cf5 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/metadata/DatanodeTable.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/metadata/DatanodeTable.java
@@ -108,6 +108,11 @@ public VALUE get(KEY key) throws RocksDatabaseException, 
CodecException {
     return table.get(key);
   }
 
+  @Override
+  public List<VALUE> multiGetSkipCache(List<KEY> keys) throws 
RocksDatabaseException, CodecException {
+    return table.multiGetSkipCache(keys);
+  }
+
   @Override
   public VALUE getIfExist(KEY key) throws RocksDatabaseException, 
CodecException {
     return table.getIfExist(key);
diff --git 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RDBTable.java
 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RDBTable.java
index d594eaf7014..afe79956600 100644
--- 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RDBTable.java
+++ 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RDBTable.java
@@ -21,12 +21,14 @@
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.List;
 import java.util.function.Supplier;
 import org.apache.hadoop.hdds.annotation.InterfaceAudience;
 import org.apache.hadoop.hdds.utils.MetadataKeyFilters.KeyPrefixFilter;
 import org.apache.hadoop.hdds.utils.db.RocksDatabase.ColumnFamily;
 import org.apache.hadoop.util.Time;
+import org.rocksdb.ByteBufferGetStatus;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -123,10 +125,32 @@ public byte[] get(byte[] key) throws 
RocksDatabaseException {
     return db.get(family, key);
   }
 
+  @Override
+  public List<byte[]> multiGetSkipCache(List<byte[]> keys) throws 
RocksDatabaseException {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    for (byte[] ignored : keys) {
+      rdbMetrics.incNumDBKeyGets();
+    }
+    return db.multiGet(family, keys);
+  }
+
   Integer get(ByteBuffer key, ByteBuffer outValue) throws 
RocksDatabaseException {
     return db.get(family, key, outValue);
   }
 
+  List<ByteBufferGetStatus> multiGetSkipCache(List<ByteBuffer> keys, 
List<ByteBuffer> values)
+      throws RocksDatabaseException {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    for (ByteBuffer ignored : keys) {
+      rdbMetrics.incNumDBKeyGets();
+    }
+    return db.multiGet(family, keys, values);
+  }
+
   /**
    * Skip checking cache and get the value mapped to the given key in byte
    * array or returns null if the key is not found.
diff --git 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RocksDatabase.java
 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RocksDatabase.java
index c2d4b8bde2b..7f4139cb773 100644
--- 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RocksDatabase.java
+++ 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/RocksDatabase.java
@@ -59,6 +59,7 @@
 import org.apache.ozone.rocksdiff.RocksDiffUtils;
 import org.apache.ratis.util.MemoizedSupplier;
 import org.apache.ratis.util.UncheckedAutoCloseable;
+import org.rocksdb.ByteBufferGetStatus;
 import org.rocksdb.ColumnFamilyDescriptor;
 import org.rocksdb.ColumnFamilyHandle;
 import org.rocksdb.Holder;
@@ -681,6 +682,21 @@ byte[] get(ColumnFamily family, byte[] key) throws 
RocksDatabaseException {
     }
   }
 
+  List<byte[]> multiGet(ColumnFamily family, List<byte[]> keys)
+      throws RocksDatabaseException {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    try (UncheckedAutoCloseable ignored = acquire()) {
+      return db.get().multiGetAsList(DEFAULT_READ_OPTION,
+          Collections.nCopies(keys.size(), family.getHandle()), keys);
+    } catch (RocksDBException e) {
+      closeOnError(e);
+      final String message = "multiGet " + keys.size() + " keys from " + 
family;
+      throw toRocksDatabaseException(this, message, e);
+    }
+  }
+
   /**
    * Get the value mapped to the given key.
    *
@@ -710,6 +726,32 @@ Integer get(ColumnFamily family, ByteBuffer key, 
ByteBuffer outValue)
     }
   }
 
+  List<ByteBufferGetStatus> multiGet(ColumnFamily family, List<ByteBuffer> 
keys, List<ByteBuffer> values)
+      throws RocksDatabaseException {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    try (UncheckedAutoCloseable ignored = acquire()) {
+      final List<ByteBufferGetStatus> statuses = 
db.get().multiGetByteBuffers(DEFAULT_READ_OPTION,
+          Collections.nCopies(keys.size(), family.getHandle()), keys, values);
+      if (LOG.isTraceEnabled()) {
+        for (ByteBufferGetStatus status: statuses) {
+          if (status.value != null) {
+            LOG.trace("multiGet: requiredSize={}, remaining={}",
+                status.requiredSize, 
status.value.asReadOnlyBuffer().remaining());
+          } else {
+            LOG.trace("multiGet: status={}", status.status);
+          }
+        }
+      }
+      return statuses;
+    } catch (RocksDBException e) {
+      closeOnError(e);
+      final String message = "multiGet " + keys.size() + " keys from " + 
family;
+      throw toRocksDatabaseException(this, message, e);
+    }
+  }
+
   public long estimateNumKeys() throws RocksDatabaseException {
     return getLongProperty(ESTIMATE_NUM_KEYS);
   }
diff --git 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/Table.java
 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/Table.java
index e740e8b5ad2..d4a03dbb0a3 100644
--- 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/Table.java
+++ 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/Table.java
@@ -18,6 +18,8 @@
 package org.apache.hadoop.hdds.utils.db;
 
 import java.io.File;
+import java.util.ArrayList;
+import java.util.Collections;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
@@ -91,6 +93,26 @@ default VALUE getSkipCache(KEY key) throws 
RocksDatabaseException, CodecExceptio
     throw new NotImplementedException("getSkipCache is not implemented");
   }
 
+  /**
+   * Skip checking cache and get the values mapped to the given keys. The
+   * returned list has the same size and order as {@code keys}; a missing key
+   * is represented by a null value.
+   *
+   * @param keys metadata keys
+   * @return values in the same order as {@code keys}
+   */
+  default List<VALUE> multiGetSkipCache(List<KEY> keys)
+      throws RocksDatabaseException, CodecException {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    List<VALUE> values = new ArrayList<>(keys.size());
+    for (KEY key : keys) {
+      values.add(getSkipCache(key));
+    }
+    return values;
+  }
+
   /**
    * Returns the value mapped to the given key in byte array or returns null
    * if the key is not found.
diff --git 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/TypedTable.java
 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/TypedTable.java
index 7a73b8938e0..b9c0c56feba 100644
--- 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/TypedTable.java
+++ 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/db/TypedTable.java
@@ -25,6 +25,7 @@
 import java.io.File;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
@@ -42,6 +43,7 @@
 import org.apache.hadoop.hdds.utils.db.cache.TableNoCache;
 import org.apache.ratis.util.Preconditions;
 import org.apache.ratis.util.function.CheckedBiFunction;
+import org.rocksdb.ByteBufferGetStatus;
 
 /**
  * Strongly typed table implementation.
@@ -228,6 +230,15 @@ public VALUE getSkipCache(KEY key) throws 
RocksDatabaseException, CodecException
     return getFromTable(key);
   }
 
+  @Override
+  public List<VALUE> multiGetSkipCache(List<KEY> keys)
+      throws RocksDatabaseException, CodecException {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    return multiGetFromTable(keys);
+  }
+
   /**
    * This method returns the value if it exists in cache, if it 
    * does not, get the value from the underlying RockDB table. If it
@@ -306,6 +317,143 @@ private VALUE getFromTable(KEY key) throws 
RocksDatabaseException, CodecExceptio
     }
   }
 
+  List<VALUE> multiGetFromTable(List<KEY> keys)
+      throws RocksDatabaseException, CodecException {
+    if (supportCodecBuffer) {
+      return multiGetFromTableWithCodecBuffer(keys);
+    }
+    return multiGetFromTableWithByteArray(keys);
+  }
+
+  List<VALUE> multiGetFromTableWithByteArray(List<KEY> keys)
+      throws RocksDatabaseException, CodecException {
+    List<byte[]> keyBytesList = new ArrayList<>(keys.size());
+    for (KEY key : keys) {
+      keyBytesList.add(encodeKey(key));
+    }
+    List<byte[]> valueBytesList = rawTable.multiGetSkipCache(keyBytesList);
+    List<VALUE> values = new ArrayList<>(keys.size());
+    for (byte[] valueBytes : valueBytesList) {
+      values.add(decodeValue(valueBytes));
+    }
+    return values;
+  }
+
+  List<VALUE> multiGetFromTableWithCodecBuffer(List<KEY> keys)
+      throws RocksDatabaseException, CodecException {
+    final int n = keys.size();
+    final List<CodecBuffer> keyBuffers = new ArrayList<>(n);
+    final List<ByteBuffer> keyByteBuffers = new ArrayList<>(n);
+    final List<CodecBuffer> valueBuffers = new ArrayList<>(n);
+    final List<ByteBuffer> valueByteBuffers = new ArrayList<>(n);
+    final int initialCapacity = bufferCapacity.get();
+    for (KEY key : keys) {
+      final CodecBuffer keyBuffer = keyCodec.toDirectCodecBuffer(key);
+      keyBuffers.add(keyBuffer);
+      keyByteBuffers.add(keyBuffer.asReadOnlyByteBuffer());
+      final CodecBuffer valueBuffer = 
CodecBuffer.allocateDirect(initialCapacity);
+      valueBuffers.add(valueBuffer);
+      addWritableValueByteBuffer(valueBuffer, initialCapacity, 
valueByteBuffers);
+    }
+
+    try {
+      final List<VALUE> values = new ArrayList<>(Collections.nCopies(n, null));
+      final List<Integer> retryIndices = new ArrayList<>();
+      final List<Integer> retrySizes = new ArrayList<>();
+      try {
+        applyMultiGetStatuses(rawTable.multiGetSkipCache(keyByteBuffers, 
valueByteBuffers),
+            valueBuffers, initialCapacity, values, retryIndices, retrySizes);
+        if (!retryIndices.isEmpty()) {
+          retryOversizedMultiGet(keyByteBuffers, retryIndices, retrySizes, 
values);
+          bufferCapacity.increase(Collections.max(retrySizes));
+        }
+        return values;
+      } finally {
+        IOUtils.closeQuietly(valueBuffers);
+      }
+    } finally {
+      IOUtils.closeQuietly(keyBuffers);
+    }
+  }
+
+  void applyMultiGetStatuses(List<ByteBufferGetStatus> statuses, 
List<CodecBuffer> valueBuffers,
+      int valueBufferCapacity, List<VALUE> values, List<Integer> retryIndices, 
List<Integer> retrySizes)
+      throws CodecException {
+    for (int i = 0; i < statuses.size(); i++) {
+      final ByteBufferGetStatus status = statuses.get(i);
+      final CodecBuffer valueBuffer = valueBuffers.get(i);
+      if (status.value == null) {
+        continue;
+      }
+      validateMultiGetStatus(status, null);
+      if (status.requiredSize > valueBufferCapacity) {
+        retryIndices.add(i);
+        retrySizes.add(status.requiredSize);
+      } else {
+        values.set(i, decodeFromMultiGetStatus(status, valueBuffer));
+      }
+    }
+  }
+
+  void retryOversizedMultiGet(List<ByteBuffer> keyByteBuffers, List<Integer> 
retryIndices,
+      List<Integer> retrySizes, List<VALUE> values) throws 
RocksDatabaseException, CodecException {
+    final int retryCount = retryIndices.size();
+    final List<ByteBuffer> retryKeyByteBuffers = new ArrayList<>(retryCount);
+    final List<CodecBuffer> retryValueBuffers = new ArrayList<>(retryCount);
+    final List<ByteBuffer> retryValueByteBuffers = new ArrayList<>(retryCount);
+    for (int j = 0; j < retryCount; j++) {
+      final int requiredSize = retrySizes.get(j);
+      
retryKeyByteBuffers.add(keyByteBuffers.get(retryIndices.get(j)).duplicate());
+      final CodecBuffer valueBuffer = CodecBuffer.allocateDirect(requiredSize);
+      retryValueBuffers.add(valueBuffer);
+      addWritableValueByteBuffer(valueBuffer, requiredSize, 
retryValueByteBuffers);
+    }
+
+    try {
+      final List<ByteBufferGetStatus> retryStatuses =
+          rawTable.multiGetSkipCache(retryKeyByteBuffers, 
retryValueByteBuffers);
+      for (int j = 0; j < retryCount; j++) {
+        final int index = retryIndices.get(j);
+        final int expectedSize = retrySizes.get(j);
+        final ByteBufferGetStatus status = retryStatuses.get(j);
+        final CodecBuffer valueBuffer = retryValueBuffers.get(j);
+        if (status.value == null) {
+          continue;
+        }
+        validateMultiGetStatus(status, expectedSize);
+        values.set(index, decodeFromMultiGetStatus(status, valueBuffer));
+      }
+    } finally {
+      IOUtils.closeQuietly(retryValueBuffers);
+    }
+  }
+
+  static void addWritableValueByteBuffer(CodecBuffer valueBuffer, int capacity,
+      List<ByteBuffer> valueByteBuffers) {
+    valueBuffer.clear();
+    final ByteBuffer out = valueBuffer.asWritableByteBuffer();
+    out.position(0).limit(capacity);
+    valueByteBuffers.add(out);
+  }
+
+  static void validateMultiGetStatus(ByteBufferGetStatus status, Integer 
expectedSize) {
+    if (status.requiredSize < 0 || status.status.getCode() != 
org.rocksdb.Status.Code.Ok) {
+      throw new IllegalStateException("status = " + status.status.getCode()
+          + "; requiredSize = " + status.requiredSize);
+    }
+    if (expectedSize != null && expectedSize.intValue() != 
status.requiredSize) {
+      throw new IllegalStateException("status = " + status.status.getCode()
+          + "; expectedSize = " + expectedSize + ", actualSize = " + 
status.requiredSize);
+    }
+  }
+
+  VALUE decodeFromMultiGetStatus(ByteBufferGetStatus status, CodecBuffer 
valueBuffer)
+      throws CodecException {
+    valueBuffer.clear();
+    valueBuffer.put(status.value.asReadOnlyBuffer());
+    return valueCodec.fromCodecBuffer(valueBuffer);
+  }
+
   /**
    * Similar to {@link #getFromTable(CodecBuffer, CodecBuffer)} except that
    * this method use {@link RDBTable#getIfExist(ByteBuffer, ByteBuffer)}.
diff --git 
a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/InMemoryTestTable.java
 
b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/InMemoryTestTable.java
index a462c9475bc..56269e62b54 100644
--- 
a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/InMemoryTestTable.java
+++ 
b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/InMemoryTestTable.java
@@ -19,6 +19,7 @@
 
 import com.google.common.primitives.UnsignedBytes;
 import java.io.File;
+import java.util.ArrayList;
 import java.util.Collections;
 import java.util.Iterator;
 import java.util.List;
@@ -93,6 +94,18 @@ public VALUE get(KEY key) {
     return map.get(key);
   }
 
+  @Override
+  public List<VALUE> multiGetSkipCache(List<KEY> keys) {
+    if (keys == null || keys.isEmpty()) {
+      return Collections.emptyList();
+    }
+    List<VALUE> values = new ArrayList<>(keys.size());
+    for (KEY key : keys) {
+      values.add(map.get(key));
+    }
+    return values;
+  }
+
   @Override
   public VALUE getIfExist(KEY key) {
     return map.get(key);
diff --git 
a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTable.java
 
b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTable.java
index dc582a537df..65a989bf8b3 100644
--- 
a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTable.java
+++ 
b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTable.java
@@ -20,6 +20,7 @@
 import static java.nio.charset.StandardCharsets.UTF_8;
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
@@ -28,6 +29,8 @@
 import static org.mockito.Mockito.when;
 
 import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.List;
 import java.util.function.Supplier;
 import org.apache.hadoop.hdds.utils.db.RocksDatabase.ColumnFamily;
 import org.junit.jupiter.api.Test;
@@ -105,5 +108,25 @@ public void testGetIfExistByteBufferFastPathReturnsValue()
     outValue.duplicate().get(readBack);
     assertArrayEquals(valueBytes, readBack);
   }
+
+  @Test
+  public void testMultiGetSkipCacheDelegatesToDatabase() throws Exception {
+    RocksDatabase db = mock(RocksDatabase.class);
+    ColumnFamily columnFamily = mock(ColumnFamily.class);
+    RDBMetrics metrics = mock(RDBMetrics.class);
+    RDBTable table = new RDBTable(db, columnFamily, metrics);
+
+    byte[] key1 = "key-1".getBytes(UTF_8);
+    byte[] key2 = "key-2".getBytes(UTF_8);
+    byte[] value1 = "value-1".getBytes(UTF_8);
+    List<byte[]> keys = Arrays.asList(key1, key2);
+    when(db.multiGet(eq(columnFamily), 
eq(keys))).thenReturn(Arrays.asList(value1, null));
+
+    List<byte[]> values = table.multiGetSkipCache(keys);
+    assertEquals(2, values.size());
+    assertArrayEquals(value1, values.get(0));
+    assertNull(values.get(1));
+    verify(db).multiGet(eq(columnFamily), eq(keys));
+  }
 }
 
diff --git 
a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTableStore.java
 
b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTableStore.java
index 586de2bf7f5..c38277e7917 100644
--- 
a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTableStore.java
+++ 
b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestRDBTableStore.java
@@ -27,13 +27,21 @@
 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.ArgumentMatchers.anyList;
+import static org.mockito.Mockito.clearInvocations;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
 
 import com.google.protobuf.ByteString;
 import java.io.File;
 import java.io.IOException;
+import java.lang.reflect.Constructor;
+import java.nio.ByteBuffer;
 import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
@@ -50,9 +58,11 @@
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
+import org.rocksdb.ByteBufferGetStatus;
 import org.rocksdb.RocksDB;
 import org.rocksdb.Statistics;
 import org.rocksdb.StatsLevel;
+import org.rocksdb.Status;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -136,6 +146,7 @@ public void tearDown() throws Exception {
       options.close();
     }
     CodecTestUtil.gc();
+    CodecBuffer.assertNoLeaks();
   }
 
   @Test
@@ -158,6 +169,94 @@ public void putGetAndEmpty() throws Exception {
     assertTrue(secondTable.isEmpty());
   }
 
+  @Test
+  public void multiGetSkipCache() throws Exception {
+    Table<byte[], byte[]> testTable = rdbStore.getTable("First");
+    byte[] key1 = 
RandomStringUtils.secure().next(10).getBytes(StandardCharsets.UTF_8);
+    byte[] key2 = 
RandomStringUtils.secure().next(10).getBytes(StandardCharsets.UTF_8);
+    byte[] missingKey = 
RandomStringUtils.secure().next(10).getBytes(StandardCharsets.UTF_8);
+    byte[] value1 = 
RandomStringUtils.secure().next(10).getBytes(StandardCharsets.UTF_8);
+    byte[] value2 = 
RandomStringUtils.secure().next(10).getBytes(StandardCharsets.UTF_8);
+    testTable.put(key1, value1);
+    testTable.put(key2, value2);
+
+    List<byte[]> values = testTable.multiGetSkipCache(Arrays.asList(key1, 
missingKey, key2));
+    assertEquals(3, values.size());
+    assertArrayEquals(value1, values.get(0));
+    assertNull(values.get(1));
+    assertArrayEquals(value2, values.get(2));
+    assertTrue(testTable.multiGetSkipCache(Collections.emptyList()).isEmpty());
+    assertTrue(testTable.multiGetSkipCache(null).isEmpty());
+  }
+
+  @Test
+  public void multiGetSkipCacheWithTypedTable() throws Exception {
+    byte[] smallKey = "small-key".getBytes(StandardCharsets.UTF_8);
+    byte[] largeKey1 = "large-key-1".getBytes(StandardCharsets.UTF_8);
+    byte[] largeKey2 = "large-key-2".getBytes(StandardCharsets.UTF_8);
+    byte[] emptyKey = "empty-key".getBytes(StandardCharsets.UTF_8);
+    byte[] smallValue = "small".getBytes(StandardCharsets.UTF_8);
+    byte[] largeValue1 = new byte[TypedTable.BUFFER_SIZE_DEFAULT + 100];
+    byte[] largeValue2 = new byte[TypedTable.BUFFER_SIZE_DEFAULT + 200];
+    Arrays.fill(largeValue1, (byte) 'a');
+    Arrays.fill(largeValue2, (byte) 'b');
+    Table<byte[], byte[]> writeTable = rdbStore.getTable("Second");
+    writeTable.put(smallKey, smallValue);
+    writeTable.put(largeKey1, largeValue1);
+    writeTable.put(largeKey2, largeValue2);
+    writeTable.put(emptyKey, new byte[0]);
+
+    Table<String, String> readTable = rdbStore.getTable("Second", 
StringCodec.get(), StringCodec.get());
+    List<String> values = readTable.multiGetSkipCache(
+        Arrays.asList("small-key", "large-key-1", "missing-key", 
"large-key-2", "empty-key"));
+    assertEquals(5, values.size());
+    assertEquals(new String(smallValue, StandardCharsets.UTF_8), 
values.get(0));
+    assertEquals(new String(largeValue1, StandardCharsets.UTF_8), 
values.get(1));
+    assertNull(values.get(2));
+    assertEquals(new String(largeValue2, StandardCharsets.UTF_8), 
values.get(3));
+    assertEquals("", values.get(4));
+  }
+
+  @Test
+  public void multiGetSkipCacheWarmsBufferCapacityHint() throws Exception {
+    byte[] largeKey = "warmup-large-key".getBytes(StandardCharsets.UTF_8);
+    byte[] largeValue = new byte[TypedTable.BUFFER_SIZE_DEFAULT + 100];
+    Arrays.fill(largeValue, (byte) 'x');
+    rdbStore.getTable("Second").put(largeKey, largeValue);
+
+    RDBTable spyTable = spy(rdbStore.getTable("Second"));
+    TypedTable<String, String> readTable = new TypedTable<>(spyTable, 
StringCodec.get(),
+        StringCodec.get(), CacheType.PARTIAL_CACHE);
+
+    readTable.multiGetSkipCache(Collections.singletonList("warmup-large-key"));
+    verify(spyTable, times(2)).multiGetSkipCache(anyList(), anyList());
+
+    clearInvocations(spyTable);
+    readTable.multiGetSkipCache(Collections.singletonList("warmup-large-key"));
+    verify(spyTable, times(1)).multiGetSkipCache(anyList(), anyList());
+  }
+
+  @Test
+  public void validateMultiGetStatusThrowsOnBadStatus() throws Exception {
+    ByteBuffer value = ByteBuffer.allocate(4);
+    ByteBufferGetStatus nonOk = newByteBufferGetStatus(
+        new Status(Status.Code.Corruption, Status.SubCode.None, "bad"), 10, 
value);
+    assertThrows(IllegalStateException.class, () -> 
TypedTable.validateMultiGetStatus(nonOk, null));
+
+    ByteBufferGetStatus negativeSize = newByteBufferGetStatus(
+        new Status(Status.Code.Incomplete, Status.SubCode.None, "too large"), 
-1, value);
+    assertThrows(IllegalStateException.class,
+        () -> TypedTable.validateMultiGetStatus(negativeSize, null));
+  }
+
+  private static ByteBufferGetStatus newByteBufferGetStatus(Status status, int 
requiredSize,
+      ByteBuffer value) throws Exception {
+    Constructor<ByteBufferGetStatus> constructor =
+        ByteBufferGetStatus.class.getDeclaredConstructor(Status.class, 
int.class, ByteBuffer.class);
+    constructor.setAccessible(true);
+    return constructor.newInstance(status, requiredSize, value);
+  }
+
   @Test
   public void delete() throws Exception {
     List<byte[]> deletedKeys = new ArrayList<>();
diff --git a/pom.xml b/pom.xml
index b6249c17e6b..64d34a5db58 100644
--- a/pom.xml
+++ b/pom.xml
@@ -1993,6 +1993,7 @@
                       <allowedImport>org.rocksdb.Env</allowedImport>
                       <allowedImport>org.rocksdb.Statistics</allowedImport>
                       <allowedImport>org.rocksdb.AbstractSlice</allowedImport>
+                      
<allowedImport>org.rocksdb.ByteBufferGetStatus</allowedImport>
                       <!-- Allow RocksDB constants and static methods to be 
used. -->
                       <allowedImport>org.rocksdb.RocksDB.*</allowedImport>
                     </allowedImports>


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to