From 3540fff7adf5212cde5c9ed0fc12cbf85e4df88b Mon Sep 17 00:00:00 2001 From: Sergey Soldatov Date: Mon, 7 Sep 2026 15:36:14 -0700 Subject: [PATCH 1/2] HDDS-16392. parse Proto2Codec values from the CodecBuffer instead of through an InputStream Co-authored-by: Claude Fable 5 --- .../apache/hadoop/hdds/utils/db/Proto2Codec.java | 9 +++------ .../apache/hadoop/hdds/utils/db/CodecTestUtil.java | 6 ++++++ .../hadoop/hdds/utils/db/Proto2CodecTestBase.java | 13 +++++++++++++ .../org/apache/hadoop/hdds/utils/db/TestCodec.java | 14 ++++++++++++++ .../ozone/om/helpers/TestTransactionInfoCodec.java | 11 +++++++++++ 5 files changed, 47 insertions(+), 6 deletions(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java index de1b8b61a5e2..6614ccf0a8c6 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java @@ -22,11 +22,9 @@ import com.google.protobuf.Parser; import jakarta.annotation.Nonnull; import java.io.IOException; -import java.io.InputStream; import java.io.OutputStream; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; -import org.apache.hadoop.hdds.utils.IOUtils; import org.apache.ratis.util.function.CheckedFunction; /** @@ -89,13 +87,12 @@ public String toString() { @Override public M fromCodecBuffer(@Nonnull CodecBuffer buffer) throws CodecException { - final InputStream in = buffer.getInputStream(); + // Parse the buffer directly, as Proto3Codec does: parsing through an InputStream makes protobuf + // allocate a 4 KB decoding buffer per call, on values which are typically a few hundred bytes. try { - return parser.parseFrom(in); + return parser.parseFrom(buffer.asReadOnlyByteBuffer()); } catch (InvalidProtocolBufferException e) { throw new CodecException("Failed to parse " + buffer + " for " + getTypeClass(), e); - } finally { - IOUtils.closeQuietly(in); } } diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java index 14ebbd63db27..3488ae30c796 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java @@ -99,6 +99,12 @@ public static void runTest(Codec codec, T original, final T fromWrappedArray = codec.fromCodecBuffer(wrapped); wrapped.release(); assertEquals(original, fromWrappedArray); + + // deserialize from direct CodecBuffer, which is what table get and iterator use + try (CodecBuffer direct = codec.toCodecBuffer(original, CodecBuffer.Allocator.getDirect())) { + assertTrue(direct.isDirect()); + assertEquals(original, codec.fromCodecBuffer(direct)); + } } public static Codec newCodecWithoutCodecBuffer(Codec codec) { diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java index 9de62ce0cc68..c4429a7495d0 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java @@ -42,6 +42,19 @@ public void testInvalidProtocolBuffer() { .contains("the input ended unexpectedly"); } + /** + * {@link Proto2Codec#fromCodecBuffer(CodecBuffer)} parses the buffer in place; a buffer which is + * not a valid message must still surface as a {@link CodecException}. + */ + @Test + public void testInvalidProtocolBufferFromCodecBuffer() { + try (CodecBuffer buffer = CodecBuffer.wrap("random".getBytes(UTF_8))) { + final CodecException exception = + assertThrows(CodecException.class, () -> getCodec().fromCodecBuffer(buffer)); + assertInstanceOf(InvalidProtocolBufferException.class, exception.getCause()); + } + } + @Test public void testFromPersistedFormat() { assertThrows(NullPointerException.class, diff --git a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java index 649c8a46f929..94c560fa6b7e 100644 --- a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java +++ b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java @@ -36,6 +36,7 @@ import java.util.UUID; import java.util.concurrent.ThreadLocalRandom; import java.util.function.Consumer; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.KeyValue; import org.apache.hadoop.hdds.utils.db.RDBBatchOperation.Bytes; import org.apache.hadoop.hdds.utils.db.managed.ManagedRocksObjectUtils; import org.junit.jupiter.api.Test; @@ -272,6 +273,19 @@ static void runTestByteStringCodec(ByteString original) throws Exception { runTest(ByteStringCodec.get(), original, original.size()); } + @Test + public void testProto2Codec() throws Exception { + final Codec codec = Proto2Codec.get(KeyValue.getDefaultInstance()); + for (int i = 0; i < NUM_LOOPS; i++) { + final KeyValue original = KeyValue.newBuilder() + .setKey("key" + i) + .setValue("value" + ThreadLocalRandom.current().nextLong()) + .build(); + runTest(codec, original, original.getSerializedSize()); + } + gc(); + } + static Executable tryCatch(Executable executable) { return tryCatch(executable, t -> LOG.info("Good!", t)); } diff --git a/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java b/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java index 7640fb20fbb8..ad2e77240028 100644 --- a/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java +++ b/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java @@ -24,6 +24,7 @@ import java.nio.charset.StandardCharsets; import org.apache.hadoop.hdds.utils.TransactionInfo; import org.apache.hadoop.hdds.utils.db.Codec; +import org.apache.hadoop.hdds.utils.db.CodecBuffer; import org.apache.hadoop.hdds.utils.db.Proto2CodecTestBase; import org.junit.jupiter.api.Test; @@ -54,4 +55,14 @@ public void testInvalidProtocolBuffer() { () -> getCodec().fromPersistedFormat("random".getBytes(StandardCharsets.UTF_8))); assertThat(ex).hasMessageContaining("Unexpected split length"); } + + @Override + @Test + public void testInvalidProtocolBufferFromCodecBuffer() { + try (CodecBuffer buffer = CodecBuffer.wrap("random".getBytes(StandardCharsets.UTF_8))) { + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, + () -> getCodec().fromCodecBuffer(buffer)); + assertThat(ex).hasMessageContaining("Unexpected split length"); + } + } } From 41c3d3c53ba3d0245dfe5dd0a6593f6eb5099111 Mon Sep 17 00:00:00 2001 From: Sergey Soldatov Date: Sat, 12 Sep 2026 12:46:19 -0700 Subject: [PATCH 2/2] Address the comments --- .../main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java index 6614ccf0a8c6..a905129c5e6f 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java @@ -87,8 +87,6 @@ public String toString() { @Override public M fromCodecBuffer(@Nonnull CodecBuffer buffer) throws CodecException { - // Parse the buffer directly, as Proto3Codec does: parsing through an InputStream makes protobuf - // allocate a 4 KB decoding buffer per call, on values which are typically a few hundred bytes. try { return parser.parseFrom(buffer.asReadOnlyByteBuffer()); } catch (InvalidProtocolBufferException e) {