This is an automated email from the ASF dual-hosted git repository.
clintropolis pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new 4eede66fcff feat: Azure SegmentRangeReader for partial loads (#20380)
4eede66fcff is described below
commit 4eede66fcff0160281f062922e972e43b1a2d40a
Author: Clint Wylie <[email protected]>
AuthorDate: Tue Sep 22 11:22:36 2026 -0700
feat: Azure SegmentRangeReader for partial loads (#20380)
---
.../storage/azure/AzureDataSegmentPuller.java | 14 ++
.../storage/azure/AzureDataSegmentPusher.java | 24 ++-
.../apache/druid/storage/azure/AzureLoadSpec.java | 45 ++++
.../storage/azure/AzureSegmentRangeReader.java | 149 +++++++++++++
.../storage/azure/AzureDataSegmentPusherTest.java | 59 ++++++
.../druid/storage/azure/AzureLoadSpecTest.java | 88 ++++++++
.../storage/azure/AzureSegmentRangeReaderTest.java | 230 +++++++++++++++++++++
7 files changed, 608 insertions(+), 1 deletion(-)
diff --git
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPuller.java
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPuller.java
index 12708297df3..65a4d9a8d6d 100644
---
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPuller.java
+++
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPuller.java
@@ -56,6 +56,20 @@ public class AzureDataSegmentPuller
this.azureAccountConfig = azureAccountConfig;
}
+ /**
+ * The storage client this puller reads through, so that {@link
AzureLoadSpec} can hand it to an
+ * {@link AzureSegmentRangeReader} rather than having the deep storage
client injected into the load spec itself.
+ */
+ AzureStorage getAzureStorage()
+ {
+ return azureStorage;
+ }
+
+ AzureAccountConfig getAccountConfig()
+ {
+ return azureAccountConfig;
+ }
+
FileUtils.FileCopyResult getSegmentFiles(
final String containerName,
final String blobPath,
diff --git
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPusher.java
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPusher.java
index 356978e8eb8..c795ecec136 100644
---
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPusher.java
+++
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentPusher.java
@@ -27,6 +27,7 @@ import org.apache.druid.guice.annotations.Global;
import org.apache.druid.java.util.common.IOE;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.java.util.common.logger.Logger;
+import org.apache.druid.segment.IndexIO;
import org.apache.druid.segment.SegmentUtils;
import org.apache.druid.segment.loading.DataSegmentPusher;
import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
@@ -187,8 +188,11 @@ public class AzureDataSegmentPusher implements
DataSegmentPusher
deleteStaleBlobs(azureDirPath, pushedBlobs);
+ // V10 unzipped is rangeable: a single druid.segment with a range-readable
header. V9 unzipped is a directory of
+ // separate smoosh files the range-read path can't consume.
+ final boolean rangeable = binaryVersion == IndexIO.V10_VERSION;
return segment.withSize(size)
- .withLoadSpec(makeLoadSpec(azureDirPath))
+ .withLoadSpec(makeLoadSpec(azureDirPath, rangeable))
.withBinaryVersion(binaryVersion);
}
@@ -280,4 +284,22 @@ public class AzureDataSegmentPusher implements
DataSegmentPusher
prefix
);
}
+
+ /**
+ * Variant that stamps {@link AzureLoadSpec#RANGEABLE} so {@link
AzureLoadSpec#openRangeReader()} can decide
+ * range-read eligibility. Used by the unzipped push path where the binary
version is known at write time.
+ */
+ private Map<String, Object> makeLoadSpec(String prefix, boolean rangeable)
+ {
+ return ImmutableMap.of(
+ "type",
+ AzureStorageDruidModule.SCHEME,
+ "containerName",
+ segmentConfig.getContainer(),
+ "blobPath",
+ prefix,
+ AzureLoadSpec.RANGEABLE,
+ rangeable
+ );
+ }
}
diff --git
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureLoadSpec.java
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureLoadSpec.java
index f6a7b347ab3..da6c4ef117c 100644
---
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureLoadSpec.java
+++
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureLoadSpec.java
@@ -21,12 +21,16 @@ package org.apache.druid.storage.azure;
import com.fasterxml.jackson.annotation.JacksonInject;
import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonTypeName;
import com.google.common.base.Preconditions;
import org.apache.druid.segment.loading.LoadSpec;
import org.apache.druid.segment.loading.SegmentLoadingException;
+import org.apache.druid.segment.loading.SegmentRangeReader;
+import org.apache.druid.utils.CompressionUtils;
+import javax.annotation.Nullable;
import java.io.File;
/**
@@ -35,6 +39,7 @@ import java.io.File;
@JsonTypeName(AzureStorageDruidModule.SCHEME)
public class AzureLoadSpec implements LoadSpec
{
+ static final String RANGEABLE = "rangeable";
@JsonProperty
private final String containerName;
@@ -42,12 +47,20 @@ public class AzureLoadSpec implements LoadSpec
@JsonProperty
private final String blobPath;
+ /**
+ * Stamped at push time when {@link AzureDataSegmentPusher} writes a segment
in a layout that supports byte-range
+ * reads. Only {@code Boolean.TRUE} enables {@link #openRangeReader()};
absence or {@code false} means full-download.
+ */
+ @Nullable
+ private final Boolean rangeable;
+
private final AzureDataSegmentPuller puller;
@JsonCreator
public AzureLoadSpec(
@JsonProperty("containerName") String containerName,
@JsonProperty("blobPath") String blobPath,
+ @JsonProperty(RANGEABLE) @Nullable Boolean rangeable,
@JacksonInject AzureDataSegmentPuller puller
)
{
@@ -55,6 +68,7 @@ public class AzureLoadSpec implements LoadSpec
Preconditions.checkNotNull(containerName);
this.containerName = containerName;
this.blobPath = blobPath;
+ this.rangeable = rangeable;
this.puller = puller;
}
@@ -63,4 +77,35 @@ public class AzureLoadSpec implements LoadSpec
{
return new LoadSpecResult(puller.getSegmentFiles(containerName, blobPath,
file).size());
}
+
+ /**
+ * Returns an {@link AzureSegmentRangeReader} when the segment was stamped
{@link #rangeable} {@code true} at push
+ * time and isn't a zip; otherwise {@code null}.
+ */
+ @Nullable
+ @Override
+ public SegmentRangeReader openRangeReader()
+ {
+ if (CompressionUtils.isZip(blobPath) || !Boolean.TRUE.equals(rangeable)) {
+ return null;
+ }
+ return new AzureSegmentRangeReader(
+ puller.getAzureStorage(),
+ containerName,
+ blobPath,
+ puller.getAccountConfig().getMaxTries()
+ );
+ }
+
+ /**
+ * Returns the range-reads-supported flag stamped at push time, or {@code
null} for legacy segments pushed before
+ * this field existed (which will load via the full-download path).
+ */
+ @JsonProperty(RANGEABLE)
+ @JsonInclude(JsonInclude.Include.NON_NULL)
+ @Nullable
+ public Boolean getRangeable()
+ {
+ return rangeable;
+ }
}
diff --git
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureSegmentRangeReader.java
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureSegmentRangeReader.java
new file mode 100644
index 00000000000..d86251205f0
--- /dev/null
+++
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureSegmentRangeReader.java
@@ -0,0 +1,149 @@
+/*
+ * 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.druid.storage.azure;
+
+import com.google.common.base.Preconditions;
+import org.apache.druid.data.input.impl.RetryingInputStream;
+import org.apache.druid.data.input.impl.prefetch.ObjectOpenFunction;
+import org.apache.druid.segment.loading.SegmentRangeReader;
+
+import javax.annotation.Nullable;
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+
+/**
+ * {@link SegmentRangeReader} backed by Azure blob range reads. The segment is
expected to be stored as raw (unzipped)
+ * blobs under a common path, i.e. the layout produced by {@code
AzureDataSegmentPusher.pushNoZip} where each segment
+ * file is uploaded as {@code blobPathPrefix + file.getName()}. Each {@link
#readRange} call resolves the target blob
+ * as {@code blobPathPrefix + filename} and opens a stream over {@code
[offset, offset + length)}.
+ * <p>
+ * The returned stream is wrapped in a {@link RetryingInputStream} with the
{@link AzureUtils#AZURE_RETRY} predicate,
+ * the same retry policy {@link AzureDataSegmentPuller} uses for full-segment
downloads. The Azure client retries a
+ * failed request on its own, but only the {@code RetryingInputStream} can
resume a read that failed partway through:
+ * it reopens at the byte offset already consumed, so a transient mid-stream
error becomes a fresh range read for the
+ * remaining bytes rather than a restart of the whole range.
+ */
+public class AzureSegmentRangeReader implements SegmentRangeReader
+{
+ private final AzureStorage azureStorage;
+ private final String containerName;
+ private final String blobPathPrefix;
+ @Nullable
+ private final Integer maxTries;
+
+ public AzureSegmentRangeReader(
+ AzureStorage azureStorage,
+ String containerName,
+ String blobPathPrefix,
+ @Nullable Integer maxTries
+ )
+ {
+ this.azureStorage = Preconditions.checkNotNull(azureStorage,
"azureStorage");
+ this.containerName = Preconditions.checkNotNull(containerName,
"containerName");
+ this.blobPathPrefix = Preconditions.checkNotNull(blobPathPrefix,
"blobPathPrefix");
+ this.maxTries = maxTries;
+ }
+
+ @Override
+ public InputStream readRange(String filename, long offset, long length)
throws IOException
+ {
+ Preconditions.checkNotNull(filename, "filename");
+ Preconditions.checkArgument(offset >= 0, "offset must be non-negative, got
[%s]", offset);
+ Preconditions.checkArgument(length >= 0, "length must be non-negative, got
[%s]", length);
+
+ if (length == 0) {
+ // SegmentFileBuilderV10 allows zero-length internal-file entries,
short-circuit
+ return new ByteArrayInputStream(new byte[0]);
+ }
+ return new RetryingInputStream<>(
+ new RangeRequest(blobPathPrefix + filename, offset, length),
+ new RangeOpenFunction(azureStorage, containerName, maxTries),
+ AzureUtils.AZURE_RETRY,
+ null
+ );
+ }
+
+ /**
+ * Immutable description of a range read. Held as the {@code object} of
{@link RetryingInputStream} so retries can
+ * reopen with knowledge of the original offset and length without
rebuilding the request from scratch.
+ */
+ private static final class RangeRequest
+ {
+ final String blobPath;
+ final long offset;
+ final long length;
+
+ RangeRequest(String blobPath, long offset, long length)
+ {
+ this.blobPath = blobPath;
+ this.offset = offset;
+ this.length = length;
+ }
+ }
+
+ /**
+ * Opens (or reopens, on retry) an Azure range read for a {@link
RangeRequest}. The {@code start} argument is the
+ * number of bytes already successfully consumed from the logical stream, so
the range still owed is
+ * {@code [request.offset + start, request.offset + request.length)}.
+ */
+ private static final class RangeOpenFunction implements
ObjectOpenFunction<RangeRequest>
+ {
+ private final AzureStorage azureStorage;
+ private final String containerName;
+ @Nullable
+ private final Integer maxTries;
+
+ RangeOpenFunction(AzureStorage azureStorage, String containerName,
@Nullable Integer maxTries)
+ {
+ this.azureStorage = azureStorage;
+ this.containerName = containerName;
+ this.maxTries = maxTries;
+ }
+
+ @Override
+ public InputStream open(RangeRequest request)
+ {
+ return open(request, 0L);
+ }
+
+ @Override
+ public InputStream open(RangeRequest request, long start)
+ {
+ final long remaining = request.length - start;
+ if (remaining <= 0) {
+ // Logically nothing left to read, only reachable if a retry fires
after the consumer drained the entire
+ // range successfully, which shouldn't happen in practice. Returning
empty keeps us robust either way.
+ return new ByteArrayInputStream(new byte[0]);
+ }
+ // A BlobStorageException is deliberately left unwrapped.
AzureUtils.AZURE_RETRY walks the cause chain but tests
+ // for IOException before BlobStorageException within each link, so
wrapping this in an IOException would make
+ // every failure look retryable and turn a 403 into ten backing-off
attempts. Raw, the predicate reads the status
+ // code and decides correctly; RetryingInputStream then hands the caller
an IOException around it either way.
+ return azureStorage.getBlockBlobInputStream(
+ request.offset + start,
+ remaining,
+ containerName,
+ request.blobPath,
+ maxTries
+ );
+ }
+ }
+}
diff --git
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPusherTest.java
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPusherTest.java
index d1eef49c08f..6bc3f316647 100644
---
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPusherTest.java
+++
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPusherTest.java
@@ -280,6 +280,65 @@ public class AzureDataSegmentPusherTest extends
EasyMockSupport
assertEquals(CONTAINER_NAME, segment.getLoadSpec().get("containerName"));
assertEquals(AzureStorageDruidModule.SCHEME,
segment.getLoadSpec().get("type"));
assertEquals(2 * DATA.length, segment.getSize());
+ // V1 (test fixture) → not V10 → rangeable stamped as false
+ assertEquals(Boolean.FALSE, segment.getLoadSpec().get("rangeable"));
+
+ verifyAll();
+ }
+
+ @Test
+ public void test_pushNoZip_v10_stampsRangeableTrue(@TempDir Path tempPath)
throws Exception
+ {
+ AzureDataSegmentPusher pusher =
+ new AzureDataSegmentPusher(azureStorage, azureAccountConfig,
segmentConfigWithPrefix, NO_ZIP_CONFIG);
+
+ // version.bin = [0, 0, 0, 0x0A] → IndexIO.V10_VERSION, a single
range-readable druid.segment
+ Files.write(new byte[]{0x0, 0x0, 0x0, 0x0A},
tempPath.resolve("version.bin").toFile());
+
+ final String expectedDir = PREFIX + "/" +
pusher.getStorageDir(SEGMENT_TO_PUSH, false);
+ azureStorage.uploadBlockBlob(
+ EasyMock.anyObject(File.class),
+ EasyMock.eq(CONTAINER_NAME),
+ EasyMock.eq(expectedDir + "/version.bin"),
+ EasyMock.eq(MAX_TRIES)
+ );
+ EasyMock.expectLastCall();
+ EasyMock.expect(azureStorage.listBlobs(CONTAINER_NAME, expectedDir + "/",
null, MAX_TRIES))
+ .andReturn(ImmutableList.of(expectedDir + "/version.bin"));
+
+ replayAll();
+
+ DataSegment segment = pusher.push(tempPath.toFile(), SEGMENT_TO_PUSH,
false);
+
+ assertEquals(10, (int) segment.getBinaryVersion());
+ assertEquals(Boolean.TRUE, segment.getLoadSpec().get("rangeable"));
+
+ verifyAll();
+ }
+
+ @Test
+ public void test_pushZip_doesNotStampRangeable(@TempDir Path tempPath)
throws Exception
+ {
+ // The zip path uses the no-flag makeLoadSpec overload; openRangeReader
returns null on the zip short circuit
+ // regardless, but the loadSpec stays compact for zipped segments by
omitting the field entirely.
+ AzureDataSegmentPusher pusher =
+ new AzureDataSegmentPusher(azureStorage, azureAccountConfig,
segmentConfigWithPrefix, ZIP_CONFIG);
+
+ Files.write(DATA, tempPath.resolve("version.bin").toFile());
+
+ azureStorage.uploadBlockBlob(
+ EasyMock.anyObject(File.class),
+ EasyMock.eq(CONTAINER_NAME),
+ EasyMock.anyString(),
+ EasyMock.eq(MAX_TRIES)
+ );
+ EasyMock.expectLastCall();
+
+ replayAll();
+
+ DataSegment segment = pusher.push(tempPath.toFile(), SEGMENT_TO_PUSH,
false);
+
+ assertFalse(segment.getLoadSpec().containsKey("rangeable"));
verifyAll();
}
diff --git
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureLoadSpecTest.java
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureLoadSpecTest.java
new file mode 100644
index 00000000000..cb56bf397fc
--- /dev/null
+++
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureLoadSpecTest.java
@@ -0,0 +1,88 @@
+/*
+ * 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.druid.storage.azure;
+
+import org.apache.druid.segment.loading.SegmentRangeReader;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.verifyNoInteractions;
+
+@ExtendWith(MockitoExtension.class)
+public class AzureLoadSpecTest
+{
+ private static final String CONTAINER = "test-container";
+ private static final String RAW_BLOB_PATH = "path/to/segment/";
+ private static final String ZIP_BLOB_PATH = "path/to/index.zip";
+
+ @Mock
+ private AzureByteSourceFactory byteSourceFactory;
+
+ @Mock
+ private AzureStorage azureStorage;
+
+ @Test
+ public void testOpenRangeReaderReturnsReaderWhenRangeableTrue()
+ {
+ final SegmentRangeReader reader = newLoadSpec(RAW_BLOB_PATH,
true).openRangeReader();
+ assertNotNull(reader);
+ assertInstanceOf(AzureSegmentRangeReader.class, reader);
+ verifyNoInteractions(azureStorage);
+ }
+
+ @Test
+ public void
testOpenRangeReaderReturnsNullForZipBlobPathEvenWhenRangeableTrue()
+ {
+ // Defensive: a zip blob can't be range-read; the zip check wins over the
flag even if hand-crafted input claims
+ // the layout is rangeable.
+ assertNull(newLoadSpec(ZIP_BLOB_PATH, true).openRangeReader());
+ verifyNoInteractions(azureStorage);
+ }
+
+ @Test
+ public void testOpenRangeReaderReturnsNullWhenRangeableFalse()
+ {
+ assertNull(newLoadSpec(RAW_BLOB_PATH, false).openRangeReader());
+ verifyNoInteractions(azureStorage);
+ }
+
+ @Test
+ public void testOpenRangeReaderReturnsNullForLegacySegmentWithoutFlag()
+ {
+ // Legacy segment (pushed before this field existed) → null flag →
full-download path.
+ assertNull(newLoadSpec(RAW_BLOB_PATH, null).openRangeReader());
+ verifyNoInteractions(azureStorage);
+ }
+
+ private AzureLoadSpec newLoadSpec(String blobPath, Boolean rangeable)
+ {
+ return new AzureLoadSpec(
+ CONTAINER,
+ blobPath,
+ rangeable,
+ new AzureDataSegmentPuller(byteSourceFactory, azureStorage, new
AzureAccountConfig())
+ );
+ }
+}
diff --git
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureSegmentRangeReaderTest.java
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureSegmentRangeReaderTest.java
new file mode 100644
index 00000000000..8ba5d5e99ab
--- /dev/null
+++
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureSegmentRangeReaderTest.java
@@ -0,0 +1,230 @@
+/*
+ * 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.druid.storage.azure;
+
+import com.azure.core.http.HttpResponse;
+import com.azure.storage.blob.models.BlobStorageException;
+import com.google.common.io.ByteStreams;
+import org.apache.druid.segment.loading.SegmentRangeReader;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.io.ByteArrayInputStream;
+import java.io.FilterInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+@ExtendWith(MockitoExtension.class)
+public class AzureSegmentRangeReaderTest
+{
+ private static final String CONTAINER = "test-container";
+ private static final String BLOB_PATH_PREFIX =
"ds/20240101T000000.000Z_20240102T000000.000Z/0/0/";
+ private static final int MAX_TRIES = 3;
+
+ @Mock
+ private AzureStorage azureStorage;
+
+ private AzureSegmentRangeReader reader;
+
+ @BeforeEach
+ public void setUp()
+ {
+ reader = new AzureSegmentRangeReader(azureStorage, CONTAINER,
BLOB_PATH_PREFIX, MAX_TRIES);
+ }
+
+ @Test
+ public void testReadRangeOpensBlobPathPrefixPlusFilenameAtOffsetAndLength()
throws IOException
+ {
+ stubRange(new byte[0]);
+
+ try (InputStream ignored = reader.readRange("druid.segment", 100, 250)) {
+ // open is performed in the RetryingInputStream constructor; the blob
read should have already happened.
+ }
+
+ verify(azureStorage).getBlockBlobInputStream(
+ 100L,
+ 250L,
+ CONTAINER,
+ BLOB_PATH_PREFIX + "druid.segment",
+ MAX_TRIES
+ );
+ }
+
+ @Test
+ public void testReadRangeBuildsDifferentBlobPathsForDifferentFilenames()
throws IOException
+ {
+ stubRange(new byte[0]);
+
+ reader.readRange("file-a", 0, 16).close();
+ reader.readRange("file-b", 0, 16).close();
+
+ final ArgumentCaptor<String> blobPathCaptor =
ArgumentCaptor.forClass(String.class);
+ verify(azureStorage, times(2)).getBlockBlobInputStream(
+ anyLong(),
+ eq(16L),
+ eq(CONTAINER),
+ blobPathCaptor.capture(),
+ eq(MAX_TRIES)
+ );
+ assertEquals(BLOB_PATH_PREFIX + "file-a",
blobPathCaptor.getAllValues().get(0));
+ assertEquals(BLOB_PATH_PREFIX + "file-b",
blobPathCaptor.getAllValues().get(1));
+ }
+
+ @Test
+ public void testReadRangeFailsFastOnNonRetryableBlobStorageException()
+ {
+ // 403 isn't retryable per AzureUtils.AZURE_RETRY, so open fails once and
RetryingInputStream surfaces it as
+ // IOException(BlobStorageException). The exception must reach the
predicate unwrapped for that to hold: AZURE_RETRY
+ // checks IOException before BlobStorageException within each link of the
cause chain, so an IOException wrapper
+ // here would report every failure as retryable and spend ten backing-off
attempts on a permission error.
+ final HttpResponse httpResponse =
org.mockito.Mockito.mock(HttpResponse.class);
+ when(httpResponse.getStatusCode()).thenReturn(403);
+ when(azureStorage.getBlockBlobInputStream(anyLong(), anyLong(),
eq(CONTAINER), eq(BLOB_PATH_PREFIX + "f"), eq(MAX_TRIES)))
+ .thenThrow(new BlobStorageException("denied", httpResponse, null));
+
+ final IOException thrown = assertThrows(IOException.class, () ->
reader.readRange("f", 0, 1));
+ assertSame(BlobStorageException.class, thrown.getCause().getClass());
+ }
+
+ @Test
+ public void testReadRangeRetriesMidStreamFromBytesAlreadyConsumed() throws
IOException
+ {
+ // First range read delivers the first 4 bytes then errors mid-stream with
a retryable IOException.
+ // RetryingInputStream should reopen for the bytes still owed, exercising
the offset math in
+ // RangeOpenFunction.open(request, start). NOTE: this test sleeps ~1s
because RetryingInputStream's first retry
+ // uses RetryUtils.BASE_SLEEP_MILLIS exponential backoff (no
@VisibleForTesting hook is accessible from here).
+ final byte[] firstChunk = {0x01, 0x02, 0x03, 0x04};
+ final byte[] secondChunk = {0x05, 0x06, 0x07, 0x08, 0x09, 0x0A};
+ final byte[] all = new byte[firstChunk.length + secondChunk.length];
+ System.arraycopy(firstChunk, 0, all, 0, firstChunk.length);
+ System.arraycopy(secondChunk, 0, all, firstChunk.length,
secondChunk.length);
+
+ when(azureStorage.getBlockBlobInputStream(anyLong(), anyLong(),
eq(CONTAINER), eq(BLOB_PATH_PREFIX + "f"), eq(MAX_TRIES)))
+ .thenReturn(failingAfter(firstChunk))
+ .thenReturn(new ByteArrayInputStream(secondChunk));
+
+ final byte[] read;
+ try (InputStream stream = reader.readRange("f", 100, 10)) {
+ read = ByteStreams.toByteArray(stream);
+ }
+ assertArrayEquals(all, read);
+
+ final ArgumentCaptor<Long> offsetCaptor =
ArgumentCaptor.forClass(Long.class);
+ final ArgumentCaptor<Long> lengthCaptor =
ArgumentCaptor.forClass(Long.class);
+ verify(azureStorage, times(2)).getBlockBlobInputStream(
+ offsetCaptor.capture(),
+ lengthCaptor.capture(),
+ eq(CONTAINER),
+ eq(BLOB_PATH_PREFIX + "f"),
+ eq(MAX_TRIES)
+ );
+ // First: the full requested range
+ assertEquals(100L, offsetCaptor.getAllValues().get(0));
+ assertEquals(10L, lengthCaptor.getAllValues().get(0));
+ // Retry: resume at (offset + bytes already consumed), asking only for
what is still owed
+ assertEquals(104L, offsetCaptor.getAllValues().get(1));
+ assertEquals(6L, lengthCaptor.getAllValues().get(1));
+ }
+
+ @Test
+ public void
testReadRangeReturnsEmptyStreamForZeroLengthWithoutContactingAzure() throws
IOException
+ {
+ // SegmentFileBuilderV10 allows zero-length internal-file entries;
readRange must accept length=0 and return an
+ // empty stream without opening a blob (a zero-length BlobRange would read
to the end of the blob instead).
+ try (InputStream stream = reader.readRange("f", 100, 0)) {
+ assertEquals(-1, stream.read());
+ }
+ verifyNoInteractions(azureStorage);
+ }
+
+ @Test
+ public void testReadRangeRejectsNegativeOffset()
+ {
+ assertThrows(IllegalArgumentException.class, () -> reader.readRange("f",
-1, 16));
+ }
+
+ @Test
+ public void testReadRangeRejectsNegativeLength()
+ {
+ assertThrows(IllegalArgumentException.class, () -> reader.readRange("f",
0, -1));
+ }
+
+ @Test
+ public void testImplementsSegmentRangeReader()
+ {
+ // Compile-time and runtime guard: ensure the implements relationship
survives refactors so callers can rely on
+ // returning AzureSegmentRangeReader from openRangeReader().
+ final SegmentRangeReader downcast = reader;
+ assertSame(reader, downcast);
+ }
+
+ private void stubRange(byte[] bytes)
+ {
+ when(azureStorage.getBlockBlobInputStream(anyLong(), anyLong(),
eq(CONTAINER), org.mockito.ArgumentMatchers.anyString(), eq(MAX_TRIES)))
+ .thenAnswer(invocation -> new ByteArrayInputStream(bytes));
+ }
+
+ /**
+ * Returns a stream that delivers the given {@code head} bytes successfully,
then raises a bare {@link IOException}
+ * (no cause) on the next read — which {@link AzureUtils#AZURE_RETRY} treats
as retryable, so
+ * {@link org.apache.druid.data.input.impl.RetryingInputStream} will close
this delegate and call the open function
+ * again with {@code start = head.length}.
+ */
+ private static InputStream failingAfter(byte[] head)
+ {
+ return new FilterInputStream(new ByteArrayInputStream(head))
+ {
+ @Override
+ public int read(byte[] b, int off, int len) throws IOException
+ {
+ final int read = super.read(b, off, len);
+ if (read < 0) {
+ throw new IOException("mid-stream failure");
+ }
+ return read;
+ }
+
+ @Override
+ public int read() throws IOException
+ {
+ final int read = super.read();
+ if (read < 0) {
+ throw new IOException("mid-stream failure");
+ }
+ return read;
+ }
+ };
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]