From 87b3f7b1e3df9115504fa0ada1fb9301e54cb727 Mon Sep 17 00:00:00 2001 From: Isabelle Date: Tue, 4 Nov 2025 14:59:03 -0800 Subject: [PATCH 1/2] all perf adjustments --- .../StructuredMessageConstants.java | 15 +++ .../StructuredMessageEncoder.java | 124 ++++++++---------- .../MessageEncoderTests.java | 62 ++++++--- 3 files changed, 113 insertions(+), 88 deletions(-) diff --git a/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageConstants.java b/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageConstants.java index 440c4521f593..859ac4c972b6 100644 --- a/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageConstants.java +++ b/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageConstants.java @@ -26,4 +26,19 @@ public final class StructuredMessageConstants { * The length of the CRC64 checksum. */ public static final int CRC64_LENGTH = 8; + + /** + * The default length of segments for version 1. + */ + public static final int V1_DEFAULT_SEGMENT_CONTENT_LENGTH = 4 * 1024 * 1024; // 4 MiB + + /** + * The maximum amount of data to encode at once. + */ + public static final int STATIC_MAXIMUM_ENCODED_DATA_LENGTH = 4 * 1024 * 1024; // 4 MiB + + /** + * The header name for the CRC64 checksum. + */ + public static final String STRUCTURED_BODY_TYPE_VALUE = "XSM/1.0; properties=crc64"; } diff --git a/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java b/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java index ca02920a13f0..d95f1d966fbc 100644 --- a/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java +++ b/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java @@ -6,12 +6,13 @@ import com.azure.core.util.logging.ClientLogger; import com.azure.storage.common.implementation.StorageCrc64Calculator; import com.azure.storage.common.implementation.StorageImplUtils; +import reactor.core.publisher.Flux; -import java.io.IOException; import java.nio.ByteBuffer; -import java.io.ByteArrayOutputStream; import java.nio.ByteOrder; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; import static com.azure.storage.common.implementation.structuredmessage.StructuredMessageConstants.CRC64_LENGTH; @@ -35,7 +36,6 @@ public class StructuredMessageEncoder { private int currentContentOffset; private int currentSegmentNumber; private int currentSegmentOffset; - private int currentMessageLength; private long messageCRC64; private final Map segmentCRC64s; @@ -66,7 +66,6 @@ public StructuredMessageEncoder(int contentLength, int segmentSize, StructuredMe this.currentSegmentOffset = 0; this.messageCRC64 = 0; this.segmentCRC64s = new HashMap<>(); - this.currentMessageLength = 0; if (numSegments > Short.MAX_VALUE) { StorageImplUtils.assertInBounds("numSegments", numSegments, 1, Short.MAX_VALUE); @@ -110,114 +109,98 @@ private byte[] generateMessageHeader() { } private byte[] generateSegmentHeader() { - int segmentHeaderSize = Math.min(segmentSize, contentLength - currentContentOffset); + int segmentContentSize = Math.min(segmentSize, contentLength - currentContentOffset); // 2 byte number, 8 byte size ByteBuffer buffer = ByteBuffer.allocate(getSegmentHeaderLength()).order(ByteOrder.LITTLE_ENDIAN); buffer.putShort((short) currentSegmentNumber); - buffer.putLong(segmentHeaderSize); + buffer.putLong(segmentContentSize); return buffer.array(); } /** - * Encodes the given buffer into a structured message format. + * Encodes the given buffer into a structured message format as a stream of ByteBuffers. * * @param unencodedBuffer The buffer to be encoded. - * @return The encoded buffer. - * @throws IOException If an error occurs while encoding the buffer. + * @return A Flux of encoded ByteBuffers. * @throws IllegalArgumentException If the buffer length exceeds the content length, or the content has already been * encoded. */ - public ByteBuffer encode(ByteBuffer unencodedBuffer) throws IOException { + public Flux encode(ByteBuffer unencodedBuffer) { StorageImplUtils.assertNotNull("unencodedBuffer", unencodedBuffer); if (currentContentOffset == contentLength) { - throw LOGGER.logExceptionAsError(new IllegalArgumentException("Content has already been encoded.")); + return Flux + .error(LOGGER.logExceptionAsError(new IllegalArgumentException("Content has already been encoded."))); } if ((unencodedBuffer.remaining() + currentContentOffset) > contentLength) { - throw LOGGER.logExceptionAsError(new IllegalArgumentException("Buffer length exceeds content length.")); + return Flux.error( + LOGGER.logExceptionAsError(new IllegalArgumentException("Buffer length exceeds content length."))); } if (!unencodedBuffer.hasRemaining()) { - return ByteBuffer.allocate(0); + return Flux.empty(); } - ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream(); + return Flux.defer(() -> { + List buffers = new ArrayList<>(); - // if we are at the beginning of the message, encode message header - if (currentMessageLength == 0) { - encodeMessageHeader(byteArrayOutputStream); - } - - while (unencodedBuffer.hasRemaining()) { - // if we are at the beginning of a segment's content, encode segment header - if (currentSegmentOffset == 0) { - encodeSegmentHeader(byteArrayOutputStream); + // if we are at the beginning of the message, encode message header + if (currentContentOffset == 0) { + buffers.add(ByteBuffer.wrap(generateMessageHeader())); } - encodeSegmentContent(unencodedBuffer, byteArrayOutputStream); - - // if we are at the end of a segment's content, encode segment footer - if (currentSegmentOffset == getSegmentContentLength()) { - encodeSegmentFooter(byteArrayOutputStream); + while (unencodedBuffer.hasRemaining()) { + // if we are at the beginning of a segment's content, encode segment header + if (currentSegmentOffset == 0) { + incrementCurrentSegment(); + buffers.add(ByteBuffer.wrap(generateSegmentHeader())); + } + + buffers.add(encodeSegmentContent(unencodedBuffer)); + + // if we are at the end of a segment's content, encode segment footer + if (currentSegmentOffset == getSegmentContentLength()) { + byte[] footer = generateSegmentFooter(); + if (footer.length > 0) { + buffers.add(ByteBuffer.wrap(footer)); + } + currentSegmentOffset = 0; + } } - } - - // if all content has been encoded, encode message footer - if (currentContentOffset == contentLength) { - encodeMessageFooter(byteArrayOutputStream); - } - - return ByteBuffer.wrap(byteArrayOutputStream.toByteArray()); - } - private void encodeMessageHeader(ByteArrayOutputStream output) { - byte[] metadata = generateMessageHeader(); - output.write(metadata, 0, metadata.length); - - currentMessageLength += metadata.length; - } - - private void encodeSegmentHeader(ByteArrayOutputStream output) { - incrementCurrentSegment(); - byte[] metadata = generateSegmentHeader(); - output.write(metadata, 0, metadata.length); + // if all content has been encoded, encode message footer + if (currentContentOffset == contentLength) { + byte[] footer = generateMessageFooter(); + if (footer.length > 0) { + buffers.add(ByteBuffer.wrap(footer)); + } + } - currentMessageLength += metadata.length; + return Flux.fromIterable(buffers); + }); } - private void encodeSegmentFooter(ByteArrayOutputStream output) { - byte[] metadata; + private byte[] generateSegmentFooter() { if (structuredMessageFlags == StructuredMessageFlags.STORAGE_CRC64) { - metadata = ByteBuffer.allocate(CRC64_LENGTH) + return ByteBuffer.allocate(CRC64_LENGTH) .order(ByteOrder.LITTLE_ENDIAN) .putLong(segmentCRC64s.get(currentSegmentNumber)) .array(); - } else { - metadata = new byte[0]; } - output.write(metadata, 0, metadata.length); - - currentMessageLength += metadata.length; - currentSegmentOffset = 0; + return new byte[0]; } - private void encodeMessageFooter(ByteArrayOutputStream output) { - byte[] metadata; + private byte[] generateMessageFooter() { if (structuredMessageFlags == StructuredMessageFlags.STORAGE_CRC64) { - metadata = ByteBuffer.allocate(CRC64_LENGTH).order(ByteOrder.LITTLE_ENDIAN).putLong(messageCRC64).array(); - } else { - metadata = new byte[0]; + return ByteBuffer.allocate(CRC64_LENGTH).order(ByteOrder.LITTLE_ENDIAN).putLong(messageCRC64).array(); } - - output.write(metadata, 0, metadata.length); - currentMessageLength += metadata.length; + return new byte[0]; } - private void encodeSegmentContent(ByteBuffer unencodedBuffer, ByteArrayOutputStream output) { + private ByteBuffer encodeSegmentContent(ByteBuffer unencodedBuffer) { int readSize = Math.min(unencodedBuffer.remaining(), getSegmentContentLength() - currentSegmentOffset); - byte[] content = new byte[readSize]; unencodedBuffer.get(content, 0, readSize); @@ -230,8 +213,7 @@ private void encodeSegmentContent(ByteBuffer unencodedBuffer, ByteArrayOutputStr currentContentOffset += readSize; currentSegmentOffset += readSize; - output.write(content, 0, content.length); - currentMessageLength += readSize; + return ByteBuffer.wrap(content); } private int calculateMessageLength() { @@ -255,7 +237,7 @@ private void incrementCurrentSegment() { * * @return The length of the message. */ - public int getMessageLength() { + public long getEncodedMessageLength() { return messageLength; } } diff --git a/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java b/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java index 6e83b7f5048c..6efbad7311f7 100644 --- a/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java +++ b/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java @@ -3,12 +3,13 @@ package com.azure.storage.common.implementation.structuredmessage; +import com.azure.core.util.FluxUtil; import com.azure.storage.common.implementation.StorageCrc64Calculator; -import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; +import reactor.core.publisher.Flux; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -18,7 +19,10 @@ import java.util.concurrent.ThreadLocalRandom; import java.util.stream.Stream; +import static com.azure.storage.common.implementation.structuredmessage.StructuredMessageConstants.V1_DEFAULT_SEGMENT_CONTENT_LENGTH; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; public class MessageEncoderTests { @@ -156,10 +160,11 @@ public void readAll(int size, int segmentSize, StructuredMessageFlags flags) thr StructuredMessageEncoder structuredMessageEncoder = new StructuredMessageEncoder(size, segmentSize, flags); - byte[] actual = structuredMessageEncoder.encode(unencodedBuffer).array(); + byte[] actual + = FluxUtil.collectBytesInByteBufferStream(structuredMessageEncoder.encode(unencodedBuffer)).block(); byte[] expected = buildStructuredMessage(unencodedBuffer, segmentSize, flags).array(); - Assertions.assertArrayEquals(expected, actual); + assertArrayEquals(expected, actual); } private static Stream readMultipleSupplier() { @@ -190,34 +195,40 @@ public void readMultiple(int segmentSize, StructuredMessageFlags flags) throws I byte[] expected = buildStructuredMessage(allWrappedData, segmentSize, flags).array(); - ByteArrayOutputStream allActualData = new ByteArrayOutputStream(); - allActualData.write(structuredMessageEncoder.encode(wrappedData1).array()); - allActualData.write(structuredMessageEncoder.encode(wrappedData2).array()); - allActualData.write(structuredMessageEncoder.encode(wrappedData3).array()); + Flux allActualFlux = structuredMessageEncoder.encode(wrappedData1) + .concatWith(structuredMessageEncoder.encode(wrappedData2)) + .concatWith(structuredMessageEncoder.encode(wrappedData3)); - Assertions.assertArrayEquals(expected, allActualData.toByteArray()); + byte[] actual = FluxUtil.collectBytesInByteBufferStream(allActualFlux).block(); + + assertArrayEquals(expected, actual); } @Test - public void emptyBuffer() throws IOException { + public void emptyBuffer() { StructuredMessageEncoder encoder = new StructuredMessageEncoder(10, 5, StructuredMessageFlags.NONE); ByteBuffer emptyBuffer = ByteBuffer.allocate(0); - ByteBuffer result = encoder.encode(emptyBuffer); - assertEquals(0, result.remaining()); + byte[] result = FluxUtil.collectBytesInByteBufferStream(encoder.encode(emptyBuffer)).block(); + assertNotNull(result); + assertEquals(0, result.length); } @Test - public void contentAlreadyEncoded() throws IOException { + public void contentAlreadyEncoded() { StructuredMessageEncoder encoder = new StructuredMessageEncoder(4, 2, StructuredMessageFlags.NONE); - encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2, 3, 4 })); - assertThrows(IllegalArgumentException.class, () -> encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2 }))); + FluxUtil.collectBytesInByteBufferStream(encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2, 3, 4 }))).block(); + assertThrows(IllegalArgumentException.class, + () -> FluxUtil.collectBytesInByteBufferStream(encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2 }))) + .block()); } @Test - public void bufferLengthExceedsContentLength() throws IOException { + public void bufferLengthExceedsContentLength() { StructuredMessageEncoder encoder = new StructuredMessageEncoder(4, 2, StructuredMessageFlags.NONE); - encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2, 3 })); - assertThrows(IllegalArgumentException.class, () -> encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2 }))); + FluxUtil.collectBytesInByteBufferStream(encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2, 3 }))).block(); + assertThrows(IllegalArgumentException.class, + () -> FluxUtil.collectBytesInByteBufferStream(encoder.encode(ByteBuffer.wrap(new byte[] { 1, 2 }))) + .block()); } @Test @@ -237,4 +248,21 @@ public void testNumSegmentsExceedsMaxValue() { assertThrows(IllegalArgumentException.class, () -> new StructuredMessageEncoder(Integer.MAX_VALUE, 1, StructuredMessageFlags.NONE)); } + + @Test + public void bigEncode() throws IOException { + byte[] data = getRandomData(262144000); + + ByteBuffer unencodedBuffer = ByteBuffer.wrap(data); + + StructuredMessageEncoder structuredMessageEncoder = new StructuredMessageEncoder(262144000, + V1_DEFAULT_SEGMENT_CONTENT_LENGTH, StructuredMessageFlags.STORAGE_CRC64); + + byte[] actual + = FluxUtil.collectBytesInByteBufferStream(structuredMessageEncoder.encode(unencodedBuffer)).block(); + byte[] expected = buildStructuredMessage(unencodedBuffer, V1_DEFAULT_SEGMENT_CONTENT_LENGTH, + StructuredMessageFlags.STORAGE_CRC64).array(); + System.out.println(expected.length); + assertArrayEquals(expected, actual); + } } From 56230719bbae380ebc22701cda1000045bbd351f Mon Sep 17 00:00:00 2001 From: Isabelle Date: Wed, 25 Feb 2026 12:38:13 -0800 Subject: [PATCH 2/2] addressing copilot comments --- .../StructuredMessageEncoder.java | 26 +++++++++---------- .../MessageEncoderTests.java | 7 +++-- 2 files changed, 18 insertions(+), 15 deletions(-) diff --git a/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java b/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java index b2ba35cb8fd3..c4951b0d8e81 100644 --- a/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java +++ b/sdk/storage/azure-storage-common/src/main/java/com/azure/storage/common/implementation/structuredmessage/StructuredMessageEncoder.java @@ -4,7 +4,6 @@ package com.azure.storage.common.implementation.structuredmessage; import com.azure.core.util.logging.ClientLogger; -import com.azure.storage.common.implementation.BufferStagingArea; import com.azure.storage.common.implementation.StorageCrc64Calculator; import com.azure.storage.common.implementation.StorageImplUtils; import reactor.core.publisher.Flux; @@ -130,21 +129,22 @@ private byte[] generateSegmentHeader() { public Flux encode(ByteBuffer unencodedBuffer) { StorageImplUtils.assertNotNull("unencodedBuffer", unencodedBuffer); - if (currentContentOffset == contentLength) { - return Flux - .error(LOGGER.logExceptionAsError(new IllegalArgumentException("Content has already been encoded."))); - } + return Flux.defer(() -> { + if (currentContentOffset == contentLength) { + return Flux.error( + LOGGER.logExceptionAsError(new IllegalArgumentException("Content has already been encoded."))); + } - if ((unencodedBuffer.remaining() + currentContentOffset) > contentLength) { - return Flux.error( - LOGGER.logExceptionAsError(new IllegalArgumentException("Buffer length exceeds content length."))); - } + if ((unencodedBuffer.remaining() + currentContentOffset) > contentLength) { + return Flux.error( + LOGGER.logExceptionAsError(new IllegalArgumentException("Buffer length exceeds content length."))); + } - if (!unencodedBuffer.hasRemaining()) { - return Flux.empty(); - } + if (!unencodedBuffer.hasRemaining()) { + return Flux.empty(); + } - return Flux.defer(() -> { + // create a list of buffers to store the encoded message List buffers = new ArrayList<>(); // if we are at the beginning of the message, encode message header diff --git a/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java b/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java index eb2e88c229fe..426ed5454abd 100644 --- a/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java +++ b/sdk/storage/azure-storage-common/src/test/java/com/azure/storage/common/implementation/structuredmessage/MessageEncoderTests.java @@ -9,6 +9,7 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; +import org.junit.jupiter.api.Disabled; import reactor.core.publisher.Flux; import java.io.ByteArrayOutputStream; @@ -16,11 +17,13 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; import java.util.Arrays; -import java.util.Objects; import java.util.concurrent.ThreadLocalRandom; import java.util.stream.Stream; import static com.azure.storage.common.implementation.structuredmessage.StructuredMessageConstants.V1_DEFAULT_SEGMENT_CONTENT_LENGTH; +import static com.azure.storage.common.implementation.structuredmessage.StructuredMessageConstants.CRC64_LENGTH; +import static com.azure.storage.common.implementation.structuredmessage.StructuredMessageConstants.V1_HEADER_LENGTH; +import static com.azure.storage.common.implementation.structuredmessage.StructuredMessageConstants.V1_SEGMENT_HEADER_LENGTH; import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; @@ -247,6 +250,7 @@ public void testNumSegmentsExceedsMaxValue() { } @Test + @Disabled("For local testing only") public void bigEncode() throws IOException { byte[] data = getRandomData(262144000); @@ -259,7 +263,6 @@ public void bigEncode() throws IOException { = FluxUtil.collectBytesInByteBufferStream(structuredMessageEncoder.encode(unencodedBuffer)).block(); byte[] expected = buildStructuredMessage(unencodedBuffer, V1_DEFAULT_SEGMENT_CONTENT_LENGTH, StructuredMessageFlags.STORAGE_CRC64).array(); - System.out.println(expected.length); assertArrayEquals(expected, actual); } }