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 a007e33a762 feat: add support to push/pull/kill segments in azure deep 
storage without zip compression (#20327)
a007e33a762 is described below

commit a007e33a762e0ed56634d4084f0263ee7967d6cb
Author: Clint Wylie <[email protected]>
AuthorDate: Fri Sep 18 11:21:38 2026 -0700

    feat: add support to push/pull/kill segments in azure deep storage without 
zip compression (#20327)
---
 docs/development/extensions-core/azure.md          |   1 +
 docs/development/extensions-core/s3.md             |   2 +-
 .../storage/azure/AzureDataSegmentKiller.java      |  11 +-
 .../storage/azure/AzureDataSegmentPuller.java      |  52 +++++-
 .../storage/azure/AzureDataSegmentPusher.java      | 132 ++++++++++++--
 .../storage/azure/AzureDataSegmentKillerTest.java  | 150 +++++++++++++---
 .../storage/azure/AzureDataSegmentPullerTest.java  |  94 +++++++++-
 .../storage/azure/AzureDataSegmentPusherTest.java  | 199 +++++++++++++++++++--
 .../druid/storage/s3/S3DataSegmentKiller.java      |   5 +-
 .../druid/storage/s3/S3DataSegmentPusher.java      |  58 +++++-
 .../storage/s3/S3DataSegmentPusherConfig.java      |   8 -
 .../java/org/apache/druid/storage/s3/S3Utils.java  |  16 +-
 .../storage/s3/S3DataSegmentPusherConfigTest.java  |   6 +-
 .../druid/storage/s3/S3DataSegmentPusherTest.java  | 167 +++++++++++++++--
 .../common/task/CompactionTaskRunBase.java         |   3 +-
 .../indexing/common/task/IngestionTestBase.java    |   5 +-
 .../AbstractParallelIndexSupervisorTaskTest.java   |   4 +-
 .../SeekableStreamIndexTaskTestBase.java           |   5 +-
 .../shuffle/ShuffleDataSegmentPusherTest.java      |   5 +-
 .../druid/msq/exec/MSQCompactionTaskRunTest.java   |   3 +-
 .../druid/msq/test/CalciteMSQTestsHelper.java      |   6 +-
 .../org/apache/druid/msq/test/MSQTestBase.java     |   3 +-
 .../segment/loading/DeepStorageSegmentConfig.java  |  73 ++++++++
 .../loading/DeepStorageSegmentConfigTest.java      | 105 +++++++++++
 .../druid/guice/LocalDataStorageDruidModule.java   |   6 +
 .../segment/loading/LocalDataSegmentPusher.java    |  11 +-
 .../loading/LocalDataSegmentPusherConfig.java      |   8 -
 .../guice/LocalDataStorageDruidModuleTest.java     |  55 ++++++
 .../loading/LocalDataSegmentPusherTest.java        |   6 +-
 29 files changed, 1075 insertions(+), 124 deletions(-)

diff --git a/docs/development/extensions-core/azure.md 
b/docs/development/extensions-core/azure.md
index b76148c7ee5..c9844f52683 100644
--- a/docs/development/extensions-core/azure.md
+++ b/docs/development/extensions-core/azure.md
@@ -60,6 +60,7 @@ Configure where to store segments using the following 
properties:
 | `druid.azure.maxTries` | Number of tries before canceling an Azure 
operation. | 3 |
 | `druid.azure.protocol` | The protocol to use to connect to the Azure Storage 
account. Either `http` or `https`. | `https` |
 | `druid.azure.storageAccountEndpointSuffix` | The Storage account endpoint to 
use. Override the default value to connect to [Azure 
Government](https://learn.microsoft.com/en-us/azure/azure-government/documentation-government-get-started-connect-to-storage#getting-started-with-storage-api)
 or storage accounts with [Azure DNS zone 
endpoints](https://learn.microsoft.com/en-us/azure/storage/common/storage-account-overview#azure-dns-zone-endpoints-preview).<br/><br/>Do
 _not_ include the stor [...]
+| `druid.storage.zip` | Whether segments are written as directories of blobs 
(`false`) or as a single `index.zip` blob (`true`). Writing segments unzipped 
requires permission to delete blobs under the segment path, so that a push 
replaces the segment already there. | `true` |
 
 #### Configure authentication
 
diff --git a/docs/development/extensions-core/s3.md 
b/docs/development/extensions-core/s3.md
index 9a00e76e4eb..43b225aabff 100644
--- a/docs/development/extensions-core/s3.md
+++ b/docs/development/extensions-core/s3.md
@@ -53,7 +53,7 @@ To use S3 for Deep Storage, you must supply [connection 
information](#configurat
 |`druid.storage.baseKey`|A prefix string that will be prepended to the object 
names for the segments published to S3 deep storage|Must be set.|
 |`druid.storage.type`|Global deep storage provider. Must be set to `s3` to 
make use of this extension.|Must be set (likely `s3`).|
 |`druid.storage.disableAcl`|Boolean flag for how object permissions are 
handled. To use ACLs, set this property to `false`. To use Object Ownership, 
set it to `true`. The permission requirements for ACLs and Object Ownership are 
different. For more information, see [S3 permissions 
settings](#s3-permissions-settings).|false|
-|`druid.storage.zip`|`true`, `false`|Whether segments in `s3` are written as 
directories (`false`) or zip files (`true`).|`false`|
+|`druid.storage.zip`|`true`, `false`|Whether segments in `s3` are written as 
directories (`false`) or zip files (`true`). Writing segments as directories 
requires `s3:DeleteObject` under the segment path, so that a push replaces the 
segment already there.|`true`|
 |`druid.storage.transfer.useTransferManager`| If true, use AWS S3 Transfer 
Manager to upload segments to S3.|true|
 |`druid.storage.transfer.minimumUploadPartSize`| Minimum size (in bytes) of 
each part in a multipart upload. Must be between 5242880 (5 MiB) and 5368709120 
(5 GiB), the range S3 permits for every part except the last. Larger parts mean 
fewer requests per upload, which matters when many tasks write concurrently 
under one key prefix.|20971520 (20 MB)|
 |`druid.storage.transfer.multipartUploadThreshold`| The file size threshold 
(in bytes) above which a file upload is converted into a multipart upload 
instead of a single PUT request.| 20971520 (20 MB)|
diff --git 
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentKiller.java
 
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentKiller.java
index ce8a6cdd388..6cee05a0934 100644
--- 
a/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentKiller.java
+++ 
b/extensions-core/azure-extensions/src/main/java/org/apache/druid/storage/azure/AzureDataSegmentKiller.java
@@ -87,7 +87,12 @@ public class AzureDataSegmentKiller implements 
DataSegmentKiller
           containerName,
           k -> new ArrayList<>()
       );
-      keysToDelete.add(blobPath);
+      if (blobPath.endsWith("/")) {
+        // segment was pushed unzipped, so the path names a directory of 
blobs; list them all to delete them
+        keysToDelete.addAll(azureStorage.listBlobs(containerName, blobPath, 
null, accountConfig.getMaxTries()));
+      } else {
+        keysToDelete.add(blobPath);
+      }
     }
 
     boolean shouldThrowException = false;
@@ -119,7 +124,9 @@ public class AzureDataSegmentKiller implements 
DataSegmentKiller
     Map<String, Object> loadSpec = segment.getLoadSpec();
     final String containerName = MapUtils.getString(loadSpec, "containerName");
     final String blobPath = MapUtils.getString(loadSpec, "blobPath");
-    final String dirPath = Paths.get(blobPath).getParent().toString();
+    final String dirPath = blobPath.endsWith("/")
+                           ? blobPath
+                           : Paths.get(blobPath).getParent() + "/";
 
     try {
       azureStorage.emptyCloudBlobDirectory(containerName, dirPath);
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 44094daade6..12708297df3 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
@@ -29,6 +29,8 @@ import org.apache.druid.utils.CompressionUtils;
 
 import java.io.File;
 import java.io.IOException;
+import java.nio.file.Paths;
+import java.util.List;
 
 /**
  * Used for Reading segment files stored in Azure based deep storage
@@ -70,18 +72,16 @@ public class AzureDataSegmentPuller
 
       final String actualBlobPath = 
AzureUtils.maybeRemoveAzurePathPrefix(blobPath, 
azureAccountConfig.getBlobStorageEndpoint());
 
-      final ByteSource byteSource = byteSourceFactory.create(containerName, 
actualBlobPath, azureStorage);
-      final FileUtils.FileCopyResult result = CompressionUtils.unzip(
-          byteSource,
-          outDir,
-          AzureUtils.AZURE_RETRY,
-          false
-      );
+      // A trailing slash means the segment was pushed unzipped 
(druid.storage.zip=false), so the path names a
+      // directory of blobs to pull individually rather than a single zip to 
unpack.
+      final FileUtils.FileCopyResult result = actualBlobPath.endsWith("/")
+                                              ? 
getSegmentFilesFromDirectory(containerName, actualBlobPath, outDir)
+                                              : 
unzipSegmentFiles(containerName, actualBlobPath, outDir);
 
       log.info("Loaded %d bytes from [%s] to [%s]", result.size(), 
actualBlobPath, outDir.getAbsolutePath());
       return result;
     }
-    catch (IOException e) {
+    catch (Exception e) {
       try {
         FileUtils.deleteDirectory(outDir);
       }
@@ -96,4 +96,40 @@ public class AzureDataSegmentPuller
       throw new SegmentLoadingException(e, e.getMessage());
     }
   }
+
+  private FileUtils.FileCopyResult unzipSegmentFiles(
+      final String containerName,
+      final String blobPath,
+      final File outDir
+  ) throws IOException
+  {
+    final ByteSource byteSource = byteSourceFactory.create(containerName, 
blobPath, azureStorage);
+    return CompressionUtils.unzip(
+        byteSource,
+        outDir,
+        AzureUtils.AZURE_RETRY,
+        false
+    );
+  }
+
+  private FileUtils.FileCopyResult getSegmentFilesFromDirectory(
+      final String containerName,
+      final String blobPathPrefix,
+      final File outDir
+  )
+  {
+    final int maxTries = azureAccountConfig.getMaxTries();
+    final List<String> blobPaths = azureStorage.listBlobs(containerName, 
blobPathPrefix, null, maxTries);
+    final FileUtils.FileCopyResult copyResult = new FileUtils.FileCopyResult();
+
+    for (final String blobPath : blobPaths) {
+      final ByteSource byteSource = byteSourceFactory.create(containerName, 
blobPath, azureStorage);
+      final File outFile = new File(outDir, 
Paths.get(blobPath).getFileName().toString());
+      copyResult.addFiles(
+          FileUtils.retryCopy(byteSource, outFile, AzureUtils.AZURE_RETRY, 
maxTries).getFiles()
+      );
+    }
+
+    return copyResult;
+  }
 }
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 000f689552f..356978e8eb8 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
@@ -24,10 +24,12 @@ import com.google.common.annotations.VisibleForTesting;
 import com.google.common.collect.ImmutableMap;
 import com.google.inject.Inject;
 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.SegmentUtils;
 import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.timeline.DataSegment;
 import org.apache.druid.utils.CompressionUtils;
 import org.joda.time.format.ISODateTimeFormat;
@@ -36,7 +38,11 @@ import java.io.File;
 import java.io.IOException;
 import java.net.URI;
 import java.nio.file.Files;
+import java.util.HashSet;
+import java.util.List;
 import java.util.Map;
+import java.util.Set;
+import java.util.stream.Collectors;
 
 /**
  * Used for writing segment files to Azure based deep storage
@@ -44,20 +50,29 @@ import java.util.Map;
 public class AzureDataSegmentPusher implements DataSegmentPusher
 {
   private static final Logger log = new Logger(AzureDataSegmentPusher.class);
+
+  /**
+   * Originally Azure always wrote segments zipped, so that is what {@code 
druid.storage.zip} falls back to.
+   */
+  private static final boolean DEFAULT_ZIP = true;
+
   private final AzureStorage azureStorage;
   private final AzureAccountConfig accountConfig;
   private final AzureDataSegmentConfig segmentConfig;
+  private final boolean zip;
 
   @Inject
   public AzureDataSegmentPusher(
       @Global AzureStorage azureStorage,
       AzureAccountConfig accountConfig,
-      AzureDataSegmentConfig segmentConfig
+      AzureDataSegmentConfig segmentConfig,
+      DeepStorageSegmentConfig deepStorageConfig
   )
   {
     this.azureStorage = azureStorage;
     this.accountConfig = accountConfig;
     this.segmentConfig = segmentConfig;
+    this.zip = deepStorageConfig.isZip(DEFAULT_ZIP);
   }
 
   @Override
@@ -87,8 +102,7 @@ public class AzureDataSegmentPusher implements 
DataSegmentPusher
       throws IOException
   {
     log.info("Uploading [%s] to Azure.", indexFilesDir);
-    final String azurePathSuffix = getAzurePath(segment, useUniquePath);
-    return pushToPath(indexFilesDir, segment, azurePathSuffix);
+    return pushToPath(indexFilesDir, segment, getStorageDir(segment, 
useUniquePath));
   }
 
   @Override
@@ -96,22 +110,39 @@ public class AzureDataSegmentPusher implements 
DataSegmentPusher
   {
     String prefix = segmentConfig.getPrefix();
     boolean prefixIsNullOrEmpty = 
org.apache.commons.lang3.StringUtils.isEmpty(prefix);
-    final String azurePath = JOINER.join(
+    final String azureBasePath = JOINER.join(
         prefixIsNullOrEmpty ? null : 
StringUtils.maybeRemoveTrailingSlash(prefix),
         storageDirSuffix
     );
 
     final int binaryVersion = SegmentUtils.getVersionFromDir(indexFilesDir);
+
+    try {
+      if (zip) {
+        return pushZip(indexFilesDir, segment, binaryVersion, azureBasePath);
+      } else {
+        return pushNoZip(indexFilesDir, segment, binaryVersion, azureBasePath);
+      }
+    }
+    catch (Exception e) {
+      throw new RuntimeException(e);
+    }
+  }
+
+  private DataSegment pushZip(
+      File indexFilesDir,
+      DataSegment segment,
+      int binaryVersion,
+      String azureBasePath
+  ) throws IOException
+  {
     File zipOutFile = null;
 
     try {
       final File outFile = zipOutFile = Files.createTempFile("index", 
".zip").toFile();
       final long size = CompressionUtils.zip(indexFilesDir, zipOutFile);
 
-      return uploadDataSegment(segment, binaryVersion, size, outFile, 
azurePath);
-    }
-    catch (Exception e) {
-      throw new RuntimeException(e);
+      return uploadDataSegment(segment, binaryVersion, size, outFile, 
getZipBlobPath(azureBasePath));
     }
     finally {
       if (zipOutFile != null) {
@@ -121,19 +152,98 @@ public class AzureDataSegmentPusher implements 
DataSegmentPusher
     }
   }
 
+  /**
+   * Uploads the segment files as they are, one blob per file, under {@code 
azureBasePath}. The resulting loadSpec
+   * blobPath is the directory itself, with a trailing slash to tell {@link 
AzureDataSegmentPuller} and
+   * {@link AzureDataSegmentKiller} that it names a directory of files rather 
than a single blob.
+   */
+  private DataSegment pushNoZip(
+      File indexFilesDir,
+      DataSegment segment,
+      int binaryVersion,
+      String azureBasePath
+  ) throws IOException
+  {
+    final File[] files = indexFilesDir.listFiles();
+    if (files == null) {
+      throw new IOE("Cannot list directory [%s]", indexFilesDir);
+    }
+
+    final String azureDirPath = azureBasePath + "/";
+    final Set<String> pushedBlobs = new HashSet<>();
+
+    long size = 0;
+    for (final File file : files) {
+      if (file.isFile()) {
+        size += file.length();
+        final String blobPath = azureDirPath + file.getName();
+        azureStorage.uploadBlockBlob(file, segmentConfig.getContainer(), 
blobPath, accountConfig.getMaxTries());
+        pushedBlobs.add(blobPath);
+      } else {
+        // Segment directories are expected to be flat.
+        throw new IOE("Unexpected subdirectory [%s]", file.getName());
+      }
+    }
+
+    deleteStaleBlobs(azureDirPath, pushedBlobs);
+
+    return segment.withSize(size)
+                  .withLoadSpec(makeLoadSpec(azureDirPath))
+                  .withBinaryVersion(binaryVersion);
+  }
+
+  /**
+   * Removes everything under {@code azureDirPath} that this push did not 
write.
+   * <p>
+   * A zipped push replaces the previous segment outright, because one {@code 
index.zip} blob overwrites another, and
+   * {@link #push} with {@code useUniquePath = false} is expected to replace a 
previous push the same way. Uploading
+   * file by file only overwrites the names the new segment happens to share, 
so without this a re-push could leave
+   * behind blobs of whatever was there before: a stale {@code index.zip} from 
a zipped push, or smoosh chunks from a
+   * larger prior v9 segment.
+   * <p>
+   * Note that this makes an unzipped push require permission to delete blobs 
under the segment path, which a zipped
+   * push does not.
+   */
+  private void deleteStaleBlobs(final String azureDirPath, final Set<String> 
pushedBlobs) throws IOException
+  {
+    final List<String> staleBlobs = azureStorage
+        .listBlobs(segmentConfig.getContainer(), azureDirPath, null, 
accountConfig.getMaxTries())
+        .stream()
+        .filter(blobPath -> !pushedBlobs.contains(blobPath))
+        .collect(Collectors.toList());
+
+    if (staleBlobs.isEmpty()) {
+      return;
+    }
+
+    log.info("Removing blobs %s left under [%s] by a previous push", 
staleBlobs, azureDirPath);
+    if (!azureStorage.batchDeleteFiles(segmentConfig.getContainer(), 
staleBlobs, accountConfig.getMaxTries())) {
+      throw new IOE(
+          "Could not remove blobs %s left under [%s] by a previous push, which 
would be loaded as part of this segment",
+          staleBlobs,
+          azureDirPath
+      );
+    }
+  }
+
   @Override
   public Map<String, Object> makeLoadSpec(URI uri)
   {
     return makeLoadSpec(uri.toString());
   }
 
+  /**
+   * Path of the {@code index.zip} blob for a segment, relative to {@link 
AzureDataSegmentConfig#getPrefix()}.
+   */
   @VisibleForTesting
   String getAzurePath(final DataSegment segment, final boolean useUniquePath)
   {
-    final String storageDir = this.getStorageDir(segment, useUniquePath);
-
-    return StringUtils.format("%s/%s", storageDir, 
AzureStorageDruidModule.INDEX_ZIP_FILE_NAME);
+    return getZipBlobPath(this.getStorageDir(segment, useUniquePath));
+  }
 
+  private static String getZipBlobPath(final String azureBasePath)
+  {
+    return StringUtils.format("%s/%s", azureBasePath, 
AzureStorageDruidModule.INDEX_ZIP_FILE_NAME);
   }
 
   @VisibleForTesting
diff --git 
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentKillerTest.java
 
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentKillerTest.java
index 40be9737d8d..82a17529d8c 100644
--- 
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentKillerTest.java
+++ 
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentKillerTest.java
@@ -29,6 +29,7 @@ import org.apache.druid.java.util.common.StringUtils;
 import org.apache.druid.segment.loading.SegmentLoadingException;
 import org.apache.druid.storage.azure.blob.CloudBlobHolder;
 import org.apache.druid.timeline.DataSegment;
+import org.apache.druid.timeline.SegmentId;
 import org.apache.druid.timeline.partition.LinearShardSpec;
 import org.easymock.Capture;
 import org.easymock.EasyMock;
@@ -44,6 +45,7 @@ import java.util.HashSet;
 import java.util.List;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
@@ -66,29 +68,31 @@ public class AzureDataSegmentKillerTest extends 
EasyMockSupport
   // BlobStorageException is not recoverable since the client attempts retries 
on it internally
   private static final Exception NON_RECOVERABLE_EXCEPTION = new 
BlobStorageException("", null, null);
 
-  private static final DataSegment DATA_SEGMENT = new DataSegment(
-      "test",
-      Intervals.of("2015-04-12/2015-04-13"),
-      "1",
-      ImmutableMap.of("containerName", CONTAINER_NAME, "blobPath", BLOB_PATH),
-      null,
-      null,
-      new LinearShardSpec(0),
-      0,
-      1
-  );
-
-  private static final DataSegment DATA_SEGMENT_2 = new DataSegment(
-      "test",
-      Intervals.of("2015-04-12/2015-04-13"),
-      "1",
-      ImmutableMap.of("containerName", CONTAINER_NAME, "blobPath", 
BLOB_PATH_2),
-      null,
-      null,
-      new LinearShardSpec(0),
-      0,
-      1
-  );
+  private static final DataSegment DATA_SEGMENT =
+      DataSegment.builder(SegmentId.of("test", 
Intervals.of("2015-04-12/2015-04-13"), "1", 0))
+                 .loadSpec(ImmutableMap.of("containerName", CONTAINER_NAME, 
"blobPath", BLOB_PATH))
+                 .shardSpec(new LinearShardSpec(0))
+                 .binaryVersion(0)
+                 .size(1)
+                 .build();
+
+  private static final DataSegment DATA_SEGMENT_2 =
+      DataSegment.builder(SegmentId.of("test", 
Intervals.of("2015-04-12/2015-04-13"), "1", 0))
+                 .loadSpec(ImmutableMap.of("containerName", CONTAINER_NAME, 
"blobPath", BLOB_PATH_2))
+                 .shardSpec(new LinearShardSpec(0))
+                 .binaryVersion(0)
+                 .size(1)
+                 .build();
+
+  // pushed with druid.storage.zip=false, so blobPath is the directory holding 
the segment files
+  private static final String UNZIPPED_BLOB_PATH = 
"test/2015-04-12T00:00:00.000Z_2015-04-13T00:00:00.000Z/1/0/";
+  private static final DataSegment UNZIPPED_DATA_SEGMENT =
+      DataSegment.builder(SegmentId.of("test", 
Intervals.of("2015-04-12/2015-04-13"), "1", 0))
+                 .loadSpec(ImmutableMap.of("containerName", CONTAINER_NAME, 
"blobPath", UNZIPPED_BLOB_PATH))
+                 .shardSpec(new LinearShardSpec(0))
+                 .binaryVersion(0)
+                 .size(1)
+                 .build();
 
   private AzureDataSegmentConfig segmentConfig;
   private AzureInputDataConfig inputDataConfig;
@@ -110,7 +114,7 @@ public class AzureDataSegmentKillerTest extends 
EasyMockSupport
   public void killTest() throws SegmentLoadingException, BlobStorageException
   {
     List<String> deletedFiles = new ArrayList<>();
-    final String dirPath = Paths.get(BLOB_PATH).getParent().toString();
+    final String dirPath = Paths.get(BLOB_PATH).getParent() + "/";
 
     EasyMock.expect(azureStorage.emptyCloudBlobDirectory(CONTAINER_NAME, 
dirPath)).andReturn(deletedFiles);
 
@@ -129,10 +133,69 @@ public class AzureDataSegmentKillerTest extends 
EasyMockSupport
     verifyAll();
   }
 
+  @Test
+  public void 
test_kill_prefixEndsAtPartitionDirectory_cannotMatchSiblingPartition()
+      throws SegmentLoadingException, BlobStorageException
+  {
+    // emptyCloudBlobDirectory matches its prefix textually rather than by 
path segment, so the prefix has to end at the
+    // partition directory boundary: without the trailing slash, killing 
partition 1 would also match partition 10.
+    final String versionDir = 
"test/2015-04-12T00:00:00.000Z_2015-04-13T00:00:00.000Z/1";
+    final String partition1 = versionDir + "/1/index.zip";
+    final String partition10 = versionDir + "/10/index.zip";
+
+    final Capture<String> prefixCapture = Capture.newInstance();
+    EasyMock.expect(azureStorage.emptyCloudBlobDirectory(
+        EasyMock.eq(CONTAINER_NAME),
+        EasyMock.capture(prefixCapture)
+    )).andReturn(ImmutableList.of(partition1));
+
+    replayAll();
+
+    final AzureDataSegmentKiller killer = new AzureDataSegmentKiller(
+        segmentConfig,
+        inputDataConfig,
+        accountConfig,
+        azureStorage,
+        azureCloudBlobIterableFactory
+    );
+
+    killer.kill(
+        DATA_SEGMENT.withLoadSpec(ImmutableMap.of("containerName", 
CONTAINER_NAME, "blobPath", partition1))
+    );
+
+    verifyAll();
+
+    assertEquals(versionDir + "/1/", prefixCapture.getValue());
+    assertFalse(partition10.startsWith(prefixCapture.getValue()));
+  }
+
+  @Test
+  public void test_kill_unzippedSegment_emptiesItsOwnDirectory() throws 
SegmentLoadingException, BlobStorageException
+  {
+    // The blobPath is already the segment's directory, so it is emptied as it 
stands. Taking its parent, as is done
+    // for a zipped segment's index.zip, would reach up into the version 
directory and take out sibling partitions.
+    EasyMock.expect(azureStorage.emptyCloudBlobDirectory(CONTAINER_NAME, 
UNZIPPED_BLOB_PATH))
+            .andReturn(ImmutableList.of(UNZIPPED_BLOB_PATH + "version.bin"));
+
+    replayAll();
+
+    final AzureDataSegmentKiller killer = new AzureDataSegmentKiller(
+        segmentConfig,
+        inputDataConfig,
+        accountConfig,
+        azureStorage,
+        azureCloudBlobIterableFactory
+    );
+
+    killer.kill(UNZIPPED_DATA_SEGMENT);
+
+    verifyAll();
+  }
+
   @Test
   public void 
test_kill_StorageExceptionExtendedErrorInformationNull_throwsException()
   {
-    String dirPath = Paths.get(BLOB_PATH).getParent().toString();
+    String dirPath = Paths.get(BLOB_PATH).getParent() + "/";
 
     EasyMock.expect(azureStorage.emptyCloudBlobDirectory(CONTAINER_NAME, 
dirPath))
             .andThrow(new BlobStorageException("", null, null));
@@ -158,7 +221,7 @@ public class AzureDataSegmentKillerTest extends 
EasyMockSupport
   @Test
   public void test_kill_runtimeException_throwsException()
   {
-    final String dirPath = Paths.get(BLOB_PATH).getParent().toString();
+    final String dirPath = Paths.get(BLOB_PATH).getParent() + "/";
 
     EasyMock.expect(azureStorage.emptyCloudBlobDirectory(CONTAINER_NAME, 
dirPath))
             .andThrow(new RuntimeException(""));
@@ -324,6 +387,39 @@ public class AzureDataSegmentKillerTest extends 
EasyMockSupport
     assertEquals(ImmutableSet.of(BLOB_PATH, BLOB_PATH_2), new 
HashSet<>(deletedFilesCapture.getValue()));
   }
 
+  @Test
+  public void test_killBatch_unzippedSegment_deletesEveryBlobInTheDirectory()
+      throws SegmentLoadingException, BlobStorageException
+  {
+    final String versionBlob = UNZIPPED_BLOB_PATH + "version.bin";
+    final String smooshBlob = UNZIPPED_BLOB_PATH + "meta.smoosh";
+
+    // an unzipped segment is a directory rather than a single blob, so its 
files have to be listed to be deleted
+    EasyMock.expect(accountConfig.getMaxTries()).andReturn(MAX_TRIES);
+    EasyMock.expect(azureStorage.listBlobs(CONTAINER_NAME, UNZIPPED_BLOB_PATH, 
null, MAX_TRIES))
+            .andReturn(ImmutableList.of(versionBlob, smooshBlob));
+
+    Capture<List<String>> deletedFilesCapture = Capture.newInstance();
+    EasyMock.expect(azureStorage.batchDeleteFiles(
+        EasyMock.eq(CONTAINER_NAME),
+        EasyMock.capture(deletedFilesCapture),
+        EasyMock.eq(null)
+    )).andReturn(true);
+
+    replayAll();
+
+    AzureDataSegmentKiller killer = new AzureDataSegmentKiller(segmentConfig, 
inputDataConfig, accountConfig, azureStorage, azureCloudBlobIterableFactory);
+
+    killer.kill(ImmutableList.of(UNZIPPED_DATA_SEGMENT, DATA_SEGMENT_2));
+
+    verifyAll();
+
+    assertEquals(
+        ImmutableSet.of(versionBlob, smooshBlob, BLOB_PATH_2),
+        new HashSet<>(deletedFilesCapture.getValue())
+    );
+  }
+
   @Test
   public void test_killBatch_runtimeException()
   {
@@ -383,7 +479,7 @@ public class AzureDataSegmentKillerTest extends 
EasyMockSupport
   public void killBatch_singleSegment() throws SegmentLoadingException, 
BlobStorageException
   {
     List<String> deletedFiles = new ArrayList<>();
-    final String dirPath = Paths.get(BLOB_PATH).getParent().toString();
+    final String dirPath = Paths.get(BLOB_PATH).getParent() + "/";
 
     // For a single segment, fall back to regular kill(DataSegment) logic
     EasyMock.expect(azureStorage.emptyCloudBlobDirectory(CONTAINER_NAME, 
dirPath)).andReturn(deletedFiles);
diff --git 
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPullerTest.java
 
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPullerTest.java
index a53e5bba405..ce19d61c566 100644
--- 
a/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPullerTest.java
+++ 
b/extensions-core/azure-extensions/src/test/java/org/apache/druid/storage/azure/AzureDataSegmentPullerTest.java
@@ -21,6 +21,7 @@ package org.apache.druid.storage.azure;
 
 import com.azure.core.http.HttpResponse;
 import com.azure.storage.blob.models.BlobStorageException;
+import com.google.common.collect.ImmutableList;
 import org.apache.druid.java.util.common.FileUtils;
 import org.apache.druid.segment.loading.SegmentLoadingException;
 import org.easymock.EasyMock;
@@ -35,6 +36,7 @@ import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
 import java.nio.file.Path;
+import java.nio.file.Paths;
 import java.util.zip.ZipEntry;
 import java.util.zip.ZipOutputStream;
 
@@ -47,6 +49,7 @@ public class AzureDataSegmentPullerTest extends 
EasyMockSupport
 {
   private static final String CONTAINER_NAME = "container";
   private static final String BLOB_PATH = "path/to/storage/index.zip";
+  private static final String UNZIPPED_BLOB_PATH = "path/to/storage/";
   private AzureStorage azureStorage;
   private AzureByteSourceFactory byteSourceFactory;
 
@@ -119,7 +122,40 @@ public class AzureDataSegmentPullerTest extends 
EasyMockSupport
   }
 
   @Test
-  public void 
test_getSegmentFiles_nonRecoverableErrorRaisedWhenPullingSegmentFiles_doNotDeleteOutputDirectory(
+  public void test_getSegmentFiles_unzippedSegment_pullsEachBlob(@TempDir Path 
sourcePath, @TempDir Path targetPath)
+      throws IOException, SegmentLoadingException
+  {
+    final AzureAccountConfig config = new AzureAccountConfig();
+    final String versionBlob = UNZIPPED_BLOB_PATH + "version.bin";
+    final String smooshBlob = UNZIPPED_BLOB_PATH + "meta.smoosh";
+    final String versionValue = "version";
+    final String smooshValue = "smoosh";
+
+    EasyMock.expect(azureStorage.listBlobs(CONTAINER_NAME, UNZIPPED_BLOB_PATH, 
null, config.getMaxTries()))
+            .andReturn(ImmutableList.of(versionBlob, smooshBlob));
+    expectBlobContents(sourcePath, versionBlob, versionValue);
+    expectBlobContents(sourcePath, smooshBlob, smooshValue);
+
+    replayAll();
+
+    AzureDataSegmentPuller puller = new 
AzureDataSegmentPuller(byteSourceFactory, azureStorage, config);
+
+    FileUtils.FileCopyResult result = puller.getSegmentFiles(CONTAINER_NAME, 
UNZIPPED_BLOB_PATH, targetPath.toFile());
+
+    // the blobs land in outDir under their own names, not unpacked from a zip
+    final File version = new File(targetPath.toFile(), "version.bin");
+    final File smoosh = new File(targetPath.toFile(), "meta.smoosh");
+    assertTrue(version.exists());
+    assertTrue(smoosh.exists());
+    assertEquals(versionValue.length(), version.length());
+    assertEquals(smooshValue.length(), smoosh.length());
+    assertEquals(versionValue.length() + smooshValue.length(), result.size());
+
+    verifyAll();
+  }
+
+  @Test
+  public void 
test_getSegmentFiles_nonRecoverableErrorRaisedWhenPullingSegmentFiles_deleteOutputDirectory(
       @TempDir Path tempPath
   )
   {
@@ -134,11 +170,47 @@ public class AzureDataSegmentPullerTest extends 
EasyMockSupport
 
     replayAll();
 
+    // an unchecked failure is still a failed load: it has to reach the caller 
as a SegmentLoadingException so that
+    // SegmentLocalCacheManager can fall back to another storage location, and 
it must not leave outDir behind
     assertThrows(
-        RuntimeException.class,
+        SegmentLoadingException.class,
         () -> puller.getSegmentFiles(CONTAINER_NAME, BLOB_PATH, 
tempPath.toFile())
     );
-    assertTrue(tempPath.toFile().exists());
+    assertFalse(tempPath.toFile().exists());
+
+    verifyAll();
+  }
+
+  @Test
+  public void 
test_getSegmentFiles_unzippedSegment_copyFails_deletesPartialOutputDirectory(
+      @TempDir Path sourcePath,
+      @TempDir Path targetPath
+  ) throws IOException
+  {
+    final AzureAccountConfig config = new AzureAccountConfig();
+    final String versionBlob = UNZIPPED_BLOB_PATH + "version.bin";
+    final String smooshBlob = UNZIPPED_BLOB_PATH + "meta.smoosh";
+
+    EasyMock.expect(azureStorage.listBlobs(CONTAINER_NAME, UNZIPPED_BLOB_PATH, 
null, config.getMaxTries()))
+            .andReturn(ImmutableList.of(versionBlob, smooshBlob));
+    // the first blob copies fine, so outDir is partially populated when the 
second one fails for good
+    expectBlobContents(sourcePath, versionBlob, "version");
+    EasyMock.expect(byteSourceFactory.create(CONTAINER_NAME, smooshBlob, 
azureStorage))
+            .andReturn(new AzureByteSource(azureStorage, CONTAINER_NAME, 
smooshBlob));
+    EasyMock.expect(azureStorage.getBlockBlobInputStream(0L, CONTAINER_NAME, 
smooshBlob))
+            .andThrow(new RuntimeException("error"));
+
+    replayAll();
+
+    AzureDataSegmentPuller puller = new 
AzureDataSegmentPuller(byteSourceFactory, azureStorage, config);
+
+    // FileUtils.retryCopy wraps the failure that exhausted its retries in a 
RuntimeException, which must not escape
+    // past the cleanup
+    assertThrows(
+        SegmentLoadingException.class,
+        () -> puller.getSegmentFiles(CONTAINER_NAME, UNZIPPED_BLOB_PATH, 
targetPath.toFile())
+    );
+    assertFalse(targetPath.toFile().exists());
 
     verifyAll();
   }
@@ -173,6 +245,22 @@ public class AzureDataSegmentPullerTest extends 
EasyMockSupport
     verifyAll();
   }
 
+  /**
+   * Sets up {@link #byteSourceFactory} and {@link #azureStorage} to serve 
{@code value} as the contents of
+   * {@code blobPath}, backed by a real file so that the copy has something to 
read.
+   */
+  private void expectBlobContents(final Path sourcePath, final String 
blobPath, final String value)
+      throws IOException
+  {
+    final File blobFile = 
Files.createFile(sourcePath.resolve(Paths.get(blobPath).getFileName().toString())).toFile();
+    Files.write(blobFile.toPath(), value.getBytes(StandardCharsets.UTF_8));
+
+    EasyMock.expect(byteSourceFactory.create(CONTAINER_NAME, blobPath, 
azureStorage))
+            .andReturn(new AzureByteSource(azureStorage, CONTAINER_NAME, 
blobPath));
+    EasyMock.expect(azureStorage.getBlockBlobInputStream(0L, CONTAINER_NAME, 
blobPath))
+            .andReturn(Files.newInputStream(blobFile.toPath()));
+  }
+
   @SuppressWarnings("SameParameterValue")
   private static File createZipTempFile(
       final Path tempPath,
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 e8595552b85..d1eef49c08f 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
@@ -20,13 +20,17 @@
 package org.apache.druid.storage.azure;
 
 import com.azure.storage.blob.models.BlobStorageException;
+import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
 import com.google.common.io.Files;
 import org.apache.druid.java.util.common.Intervals;
 import org.apache.druid.java.util.common.MapUtils;
 import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.timeline.DataSegment;
 import org.apache.druid.timeline.partition.LinearShardSpec;
+import org.easymock.Capture;
 import org.easymock.EasyMock;
 import org.easymock.EasyMockSupport;
 import org.junit.jupiter.api.BeforeEach;
@@ -63,6 +67,8 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
       0,
       1
   );
+  private static final DeepStorageSegmentConfig ZIP_CONFIG = new 
DeepStorageSegmentConfig(true);
+  private static final DeepStorageSegmentConfig NO_ZIP_CONFIG = new 
DeepStorageSegmentConfig(false);
   private static final byte[] DATA = new byte[]{0x0, 0x0, 0x0, 0x1};
   private static final String UNIQUE_MATCHER_NO_PREFIX = 
"foo/20150101T000000\\.000Z_20160101T000000\\.000Z/0/0/[A-Za-z0-9-]{36}/index\\.zip";
   private static final String UNIQUE_MATCHER_PREFIX = PREFIX + "/" + 
UNIQUE_MATCHER_NO_PREFIX;
@@ -107,8 +113,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   public void test_push_nonUniquePathNoPrefix_succeeds(@TempDir Path tempPath) 
throws Exception
   {
     boolean useUniquePath = false;
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithoutPrefix
-    );
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithoutPrefix, ZIP_CONFIG);
 
     // Create a mock segment on disk
     File tmp = tempPath.resolve("version.bin").toFile();
@@ -138,8 +143,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   public void test_push_nonUniquePathWithPrefix_succeeds(@TempDir Path 
tempPath) throws Exception
   {
     boolean useUniquePath = false;
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix
-    );
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix, ZIP_CONFIG);
 
     // Create a mock segment on disk
     File tmp = tempPath.resolve("version.bin").toFile();
@@ -170,7 +174,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   public void test_push_uniquePathNoPrefix_succeeds(@TempDir Path tempPath) 
throws Exception
   {
     boolean useUniquePath = true;
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithoutPrefix);
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithoutPrefix, ZIP_CONFIG);
 
     // Create a mock segment on disk
     File tmp = tempPath.resolve("version.bin").toFile();
@@ -205,7 +209,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   public void test_push_uniquePath_succeeds(@TempDir Path tempPath) throws 
Exception
   {
     boolean useUniquePath = true;
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix);
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix, ZIP_CONFIG);
 
     // Create a mock segment on disk
     File tmp = tempPath.resolve("version.bin").toFile();
@@ -237,11 +241,186 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
     verifyAll();
   }
 
+  @Test
+  public void test_pushNoZip_uploadsEachFile_succeeds(@TempDir Path tempPath) 
throws Exception
+  {
+    AzureDataSegmentPusher pusher =
+        new AzureDataSegmentPusher(azureStorage, azureAccountConfig, 
segmentConfigWithPrefix, NO_ZIP_CONFIG);
+
+    // Create a mock segment on disk, of more than one file, since the point 
of not zipping is to keep them separate
+    Files.write(DATA, tempPath.resolve("version.bin").toFile());
+    Files.write(DATA, tempPath.resolve("meta.smoosh").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();
+    azureStorage.uploadBlockBlob(
+        EasyMock.anyObject(File.class),
+        EasyMock.eq(CONTAINER_NAME),
+        EasyMock.eq(expectedDir + "/meta.smoosh"),
+        EasyMock.eq(MAX_TRIES)
+    );
+    EasyMock.expectLastCall();
+
+    // nothing else is under the path, so there is nothing to clean up
+    EasyMock.expect(azureStorage.listBlobs(CONTAINER_NAME, expectedDir + "/", 
null, MAX_TRIES))
+            .andReturn(ImmutableList.of(expectedDir + "/version.bin", 
expectedDir + "/meta.smoosh"));
+
+    replayAll();
+
+    DataSegment segment = pusher.push(tempPath.toFile(), SEGMENT_TO_PUSH, 
false);
+
+    // the trailing slash is what marks the blobPath as a directory of files 
rather than a single blob
+    assertEquals(expectedDir + "/", segment.getLoadSpec().get("blobPath"));
+    assertEquals(CONTAINER_NAME, segment.getLoadSpec().get("containerName"));
+    assertEquals(AzureStorageDruidModule.SCHEME, 
segment.getLoadSpec().get("type"));
+    assertEquals(2 * DATA.length, segment.getSize());
+
+    verifyAll();
+  }
+
+  @Test
+  public void test_pushNoZip_removesBlobsLeftByPreviousPush(@TempDir Path 
tempPath) throws Exception
+  {
+    AzureDataSegmentPusher pusher =
+        new AzureDataSegmentPusher(azureStorage, azureAccountConfig, 
segmentConfigWithPrefix, NO_ZIP_CONFIG);
+
+    Files.write(DATA, 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();
+
+    // an index.zip from a zipped push and a smoosh chunk from a larger prior 
segment: nothing in the new segment
+    // references either, but the puller would list and download both
+    final String staleZip = expectedDir + "/index.zip";
+    final String staleChunk = expectedDir + "/00001.smoosh";
+    EasyMock.expect(azureStorage.listBlobs(CONTAINER_NAME, expectedDir + "/", 
null, MAX_TRIES))
+            .andReturn(ImmutableList.of(expectedDir + "/version.bin", 
staleZip, staleChunk));
+
+    final Capture<Iterable<String>> deleted = Capture.newInstance();
+    EasyMock.expect(azureStorage.batchDeleteFiles(
+        EasyMock.eq(CONTAINER_NAME),
+        EasyMock.capture(deleted),
+        EasyMock.eq(MAX_TRIES)
+    )).andReturn(true);
+
+    replayAll();
+
+    pusher.push(tempPath.toFile(), SEGMENT_TO_PUSH, false);
+
+    verifyAll();
+
+    // only the blobs this push did not write are removed
+    assertEquals(ImmutableSet.of(staleZip, staleChunk), 
ImmutableSet.copyOf(deleted.getValue()));
+  }
+
+  @Test
+  public void test_pushNoZip_staleBlobDeleteFails_throwsException(@TempDir 
Path tempPath) throws Exception
+  {
+    AzureDataSegmentPusher pusher =
+        new AzureDataSegmentPusher(azureStorage, azureAccountConfig, 
segmentConfigWithPrefix, NO_ZIP_CONFIG);
+
+    Files.write(DATA, 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 + "/index.zip"));
+    EasyMock.expect(azureStorage.batchDeleteFiles(
+        EasyMock.eq(CONTAINER_NAME),
+        EasyMock.anyObject(),
+        EasyMock.eq(MAX_TRIES)
+    )).andReturn(false);
+
+    replayAll();
+
+    // failing the push is better than publishing a segment that would load 
the leftovers alongside its own files
+    assertThrows(
+        RuntimeException.class,
+        () -> pusher.push(tempPath.toFile(), SEGMENT_TO_PUSH, false)
+    );
+
+    verifyAll();
+  }
+
+  @Test
+  public void test_pushNoZip_subdirectory_throwsException(@TempDir Path 
tempPath) throws Exception
+  {
+    AzureDataSegmentPusher pusher =
+        new AzureDataSegmentPusher(azureStorage, azureAccountConfig, 
segmentConfigWithPrefix, NO_ZIP_CONFIG);
+
+    Files.write(DATA, tempPath.resolve("version.bin").toFile());
+    // Segment directories are expected to be flat.
+    java.nio.file.Files.createDirectory(tempPath.resolve("nested"));
+
+    // version.bin may be uploaded first or not at all, depending on the order 
the directory lists in
+    azureStorage.uploadBlockBlob(
+        EasyMock.anyObject(File.class),
+        EasyMock.anyString(),
+        EasyMock.anyString(),
+        EasyMock.anyInt()
+    );
+    EasyMock.expectLastCall().anyTimes();
+
+    replayAll();
+
+    assertThrows(
+        RuntimeException.class,
+        () -> pusher.push(tempPath.toFile(), SEGMENT_TO_PUSH, false)
+    );
+
+    verifyAll();
+  }
+
+  @Test
+  public void test_pushToPath_nonSegmentDirSuffix_appendsIndexZip(@TempDir 
Path tempPath) throws Exception
+  {
+    // The shuffle intermediary data manager pushes to a path of its own 
choosing, which is a directory like any other
+    AzureDataSegmentPusher pusher =
+        new AzureDataSegmentPusher(azureStorage, azureAccountConfig, 
segmentConfigWithoutPrefix, ZIP_CONFIG);
+
+    Files.write(DATA, tempPath.resolve("version.bin").toFile());
+
+    azureStorage.uploadBlockBlob(
+        EasyMock.anyObject(File.class),
+        EasyMock.eq(CONTAINER_NAME),
+        EasyMock.eq("shuffle-data/supervisorId/partition/index.zip"),
+        EasyMock.eq(MAX_TRIES)
+    );
+    EasyMock.expectLastCall();
+
+    replayAll();
+
+    DataSegment segment = pusher.pushToPath(tempPath.toFile(), 
SEGMENT_TO_PUSH, "shuffle-data/supervisorId/partition");
+
+    assertEquals("shuffle-data/supervisorId/partition/index.zip", 
segment.getLoadSpec().get("blobPath"));
+
+    verifyAll();
+  }
+
   @Test
   public void test_push_exception_throwsException(@TempDir Path tempPath) 
throws Exception
   {
     boolean useUniquePath = true;
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix);
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix, ZIP_CONFIG);
 
     // Create a mock segment on disk
     File tmp = tempPath.resolve("version.bin").toFile();
@@ -264,7 +443,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   @Test
   public void getAzurePathsTest()
   {
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix);
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix, ZIP_CONFIG);
     final String storageDir = pusher.getStorageDir(DATA_SEGMENT, false);
     final String azurePath = pusher.getAzurePath(DATA_SEGMENT, false);
 
@@ -277,7 +456,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   @Test
   public void uploadDataSegmentTest() throws BlobStorageException, IOException
   {
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix);
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix, ZIP_CONFIG);
     final int binaryVersion = 9;
     final File compressedSegmentData = new File("index.zip");
     final String azurePath = pusher.getAzurePath(DATA_SEGMENT, false);
@@ -307,7 +486,7 @@ public class AzureDataSegmentPusherTest extends 
EasyMockSupport
   @Test
   public void storageDirContainsNoColonsTest()
   {
-    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix);
+    AzureDataSegmentPusher pusher = new AzureDataSegmentPusher(azureStorage, 
azureAccountConfig, segmentConfigWithPrefix, ZIP_CONFIG);
     DataSegment withColons = 
DATA_SEGMENT.withVersion("2018-01-05T14:54:09.295Z");
     String segmentPath = pusher.getStorageDir(withColons, false);
     assertFalse(segmentPath.contains(":"), "Path should not contain any 
columns");
diff --git 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentKiller.java
 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentKiller.java
index a21080b4da7..20cf215bb66 100644
--- 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentKiller.java
+++ 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentKiller.java
@@ -49,9 +49,6 @@ public class S3DataSegmentKiller implements DataSegmentKiller
 {
   private static final Logger log = new Logger(S3DataSegmentKiller.class);
 
-  // AWS has max limit of 1000 objects that can be requested to be deleted at 
a time.
-  private static final int MAX_MULTI_OBJECT_DELETE_SIZE = 1000;
-
   private static final int NUM_RETRIES = 3;
 
   /**
@@ -152,7 +149,7 @@ public class S3DataSegmentKiller implements 
DataSegmentKiller
   )
   {
     boolean hadException = false;
-    for (List<ObjectIdentifier> chunkOfKeys : Lists.partition(keysToDelete, 
MAX_MULTI_OBJECT_DELETE_SIZE)) {
+    for (List<ObjectIdentifier> chunkOfKeys : Lists.partition(keysToDelete, 
S3Utils.MAX_MULTI_OBJECT_DELETE_SIZE)) {
       try {
         log.info("Deleting %d segment files from S3 bucket[%s]", 
chunkOfKeys.size(), s3Bucket);
         S3Utils.deleteBucketKeys(s3Client, s3Bucket, chunkOfKeys, NUM_RETRIES);
diff --git 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusher.java
 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusher.java
index b7da01bfbaa..288802bbc5a 100644
--- 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusher.java
+++ 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusher.java
@@ -27,6 +27,7 @@ import org.apache.druid.java.util.emitter.EmittingLogger;
 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;
 import org.apache.druid.timeline.DataSegment;
 import org.apache.druid.utils.CompressionUtils;
 import software.amazon.awssdk.services.s3.model.S3Exception;
@@ -35,23 +36,33 @@ import java.io.File;
 import java.io.IOException;
 import java.net.URI;
 import java.nio.file.Files;
+import java.util.HashSet;
 import java.util.Map;
+import java.util.Set;
 
 public class S3DataSegmentPusher implements DataSegmentPusher
 {
   private static final EmittingLogger log = new 
EmittingLogger(S3DataSegmentPusher.class);
 
+  /**
+   * Originally S3 always wrote segments zipped, so that is what {@code 
druid.storage.zip} falls back to.
+   */
+  private static final boolean DEFAULT_ZIP = true;
+
   private final ServerSideEncryptingAmazonS3 s3Client;
   private final S3DataSegmentPusherConfig config;
+  private final boolean zip;
 
   @Inject
   public S3DataSegmentPusher(
       ServerSideEncryptingAmazonS3 s3Client,
-      S3DataSegmentPusherConfig config
+      S3DataSegmentPusherConfig config,
+      DeepStorageSegmentConfig deepStorageConfig
   )
   {
     this.s3Client = s3Client;
     this.config = config;
+    this.zip = deepStorageConfig.isZip(DEFAULT_ZIP);
   }
 
   @Override
@@ -64,7 +75,7 @@ public class S3DataSegmentPusher implements DataSegmentPusher
   @Override
   public DataSegment pushToPath(File indexFilesDir, DataSegment inSegment, 
String storageDirSuffix) throws IOException
   {
-    if (config.isZip()) {
+    if (zip) {
       final String s3Path = S3Utils.constructSegmentPath(config.getBaseKey(), 
storageDirSuffix);
       log.debug("Copying segment[%s] to S3 at location[%s]", 
inSegment.getId(), s3Path);
       return pushZip(indexFilesDir, inSegment, s3Path);
@@ -114,15 +125,18 @@ public class S3DataSegmentPusher implements 
DataSegmentPusher
       throw new IOE("Cannot list directory [%s]", indexFilesDir);
     }
 
+    final Set<String> pushedKeys = new HashSet<>();
+
     long size = 0;
     for (final File file : files) {
       if (file.isFile()) {
         size += file.length();
+        final String key = s3Path + file.getName();
 
         try {
           S3Utils.retryS3Operation(
               () -> {
-                S3Utils.uploadFileIfPossible(s3Client, config.getDisableAcl(), 
config.getBucket(), s3Path + file.getName(), file);
+                S3Utils.uploadFileIfPossible(s3Client, config.getDisableAcl(), 
config.getBucket(), key, file);
                 return null;
               }
           );
@@ -133,12 +147,16 @@ public class S3DataSegmentPusher implements 
DataSegmentPusher
         catch (Exception e) {
           throw new RuntimeException(e);
         }
+
+        pushedKeys.add(key);
       } else {
         // Segment directories are expected to be flat.
         throw new IOE("Unexpected subdirectory [%s]", file.getName());
       }
     }
 
+    deleteStaleObjects(s3Path, pushedKeys);
+
     final int binaryVersion = SegmentUtils.getVersionFromDir(indexFilesDir);
     // 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.
@@ -148,6 +166,40 @@ public class S3DataSegmentPusher implements 
DataSegmentPusher
                       .withBinaryVersion(binaryVersion);
   }
 
+  /**
+   * Removes everything under {@code s3Path} that this push did not write.
+   * <p>
+   * A zipped push replaces the previous segment outright, because one {@code 
index.zip} object overwrites another, and
+   * {@link #push} with {@code useUniquePath = false} is expected to replace a 
previous push the same way. Uploading
+   * file by file only overwrites the names the new segment happens to share, 
so without this a re-push could leave
+   * behind objects of whatever was there before: a stale {@code index.zip} 
from a zipped push, or smoosh chunks from a
+   * larger prior v9 segment.
+   * <p>
+   * Note that this makes an unzipped push require permission to delete 
objects under the segment path, which a zipped
+   * push does not.
+   */
+  private void deleteStaleObjects(final String s3Path, final Set<String> 
pushedKeys) throws IOException
+  {
+    try {
+      S3Utils.deleteObjectsInPath(
+          s3Client,
+          config.getMaxListingLength(),
+          config.getBucket(),
+          s3Path,
+          object -> !pushedKeys.contains(object.key())
+      );
+    }
+    catch (Exception e) {
+      throw new IOE(
+          e,
+          "Could not remove objects left under [s3://%s/%s] by a previous 
push, which would be loaded as part of this"
+          + " segment",
+          config.getBucket(),
+          s3Path
+      );
+    }
+  }
+
   @Override
   public Map<String, Object> makeLoadSpec(URI finalIndexZipFilePath)
   {
diff --git 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfig.java
 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfig.java
index 03592cca3c8..94209c68276 100644
--- 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfig.java
+++ 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfig.java
@@ -39,9 +39,6 @@ public class S3DataSegmentPusherConfig
   @Min(1)
   private int maxListingLength = 1024;
 
-  @JsonProperty
-  private boolean zip = true;
-
   public void setBucket(String bucket)
   {
     this.bucket = bucket;
@@ -81,9 +78,4 @@ public class S3DataSegmentPusherConfig
   {
     return maxListingLength;
   }
-
-  public boolean isZip()
-  {
-    return zip;
-  }
 }
diff --git 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java
 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java
index 0ad2c17cd50..bef5b6ecc01 100644
--- 
a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java
+++ 
b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java
@@ -74,6 +74,13 @@ public class S3Utils
    */
   public static final String ERROR_ENTITY_TOO_LARGE = "EntityTooLarge";
 
+  /**
+   * Most keys S3 accepts in a single {@code DeleteObjects} request:
+   * <a 
href="https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteObjects.html";>DeleteObjects</a>.
 Unlike
+   * {@code max-keys} on a listing, which S3 silently caps, exceeding this 
fails the request.
+   */
+  public static final int MAX_MULTI_OBJECT_DELETE_SIZE = 1000;
+
   public static final Predicate<Throwable> S3RETRY = new Predicate<>()
   {
     @Override
@@ -360,7 +367,12 @@ public class S3Utils
   {
     log.debug("Deleting directory at bucket: [%s], path: [%s]", bucket, 
prefix);
 
-    final List<ObjectIdentifier> keysToDelete = new 
ArrayList<>(maxListingLength);
+    // A listing page and a multi-object delete do not have the same ceiling: 
S3 quietly returns only the first 1000
+    // keys when max-keys is larger, but rejects a DeleteObjects request 
carrying more than 1000 of them. Callers pass
+    // maxListingLength for the listing, so bound the delete batches 
separately instead of assuming the caller's
+    // listing size is small enough to be one.
+    final int deleteBatchSize = Math.min(maxListingLength, 
MAX_MULTI_OBJECT_DELETE_SIZE);
+    final List<ObjectIdentifier> keysToDelete = new 
ArrayList<>(deleteBatchSize);
     final Iterator<S3ObjectWithBucket> iterator = objectSummaryIterator(
         s3Client,
         ImmutableList.of(new CloudObjectLocation(bucket, prefix).toUri("s3")),
@@ -371,7 +383,7 @@ public class S3Utils
       final S3ObjectWithBucket nextObject = iterator.next();
       if (filter.apply(nextObject.getS3Object())) {
         
keysToDelete.add(ObjectIdentifier.builder().key(nextObject.getKey()).build());
-        if (keysToDelete.size() == maxListingLength) {
+        if (keysToDelete.size() == deleteBatchSize) {
           deleteBucketKeys(s3Client, bucket, keysToDelete, maxRetries);
           keysToDelete.clear();
         }
diff --git 
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfigTest.java
 
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfigTest.java
index 9a4e58bd7a6..0419d640788 100644
--- 
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfigTest.java
+++ 
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherConfigTest.java
@@ -40,7 +40,7 @@ public class S3DataSegmentPusherConfigTest
   public void testSerialization() throws IOException
   {
     String jsonConfig = "{\"bucket\":\"bucket1\",\"baseKey\":\"dataSource1\","
-                        + 
"\"disableAcl\":false,\"maxListingLength\":2000,\"zip\":true}";
+                        + "\"disableAcl\":false,\"maxListingLength\":2000}";
 
     S3DataSegmentPusherConfig config = JSON_MAPPER.readValue(jsonConfig, 
S3DataSegmentPusherConfig.class);
     Map<String, String> expected = JSON_MAPPER.readValue(jsonConfig, 
Map.class);
@@ -53,7 +53,7 @@ public class S3DataSegmentPusherConfigTest
   {
     String jsonConfig = "{\"bucket\":\"bucket1\",\"baseKey\":\"dataSource1\"}";
     String expectedJsonConfig = 
"{\"bucket\":\"bucket1\",\"baseKey\":\"dataSource1\","
-                                + 
"\"disableAcl\":false,\"maxListingLength\":1024,\"zip\":true}";
+                                + 
"\"disableAcl\":false,\"maxListingLength\":1024}";
     S3DataSegmentPusherConfig config = JSON_MAPPER.readValue(jsonConfig, 
S3DataSegmentPusherConfig.class);
     Map<String, String> expected = JSON_MAPPER.readValue(expectedJsonConfig, 
Map.class);
     Map<String, String> actual = 
JSON_MAPPER.readValue(JSON_MAPPER.writeValueAsString(config), Map.class);
@@ -64,7 +64,7 @@ public class S3DataSegmentPusherConfigTest
   public void testSerializationValidatingMaxListingLength() throws IOException
   {
     String jsonConfig = "{\"bucket\":\"bucket1\",\"baseKey\":\"dataSource1\","
-                        + 
"\"disableAcl\":false,\"maxListingLength\":-1,\"zip\":true}";
+                        + "\"disableAcl\":false,\"maxListingLength\":-1}";
     Validator validator = 
Validation.buildDefaultValidatorFactory().getValidator();
 
     S3DataSegmentPusherConfig config = JSON_MAPPER.readValue(jsonConfig, 
S3DataSegmentPusherConfig.class);
diff --git 
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherTest.java
 
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherTest.java
index a7a31310309..5e89f5a3897 100644
--- 
a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherTest.java
+++ 
b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3DataSegmentPusherTest.java
@@ -19,27 +19,40 @@
 
 package org.apache.druid.storage.s3;
 
+import com.google.common.collect.ImmutableSet;
 import com.google.common.io.Files;
 import org.apache.druid.error.DruidException;
 import org.apache.druid.java.util.common.Intervals;
+import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.timeline.DataSegment;
 import org.apache.druid.timeline.partition.NoneShardSpec;
+import org.easymock.Capture;
+import org.easymock.CaptureType;
 import org.easymock.EasyMock;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 import software.amazon.awssdk.awscore.exception.AwsErrorDetails;
+import software.amazon.awssdk.services.s3.model.DeleteObjectsRequest;
+import software.amazon.awssdk.services.s3.model.DeleteObjectsResponse;
 import software.amazon.awssdk.services.s3.model.Grant;
 import software.amazon.awssdk.services.s3.model.Grantee;
+import software.amazon.awssdk.services.s3.model.ListObjectsV2Request;
+import software.amazon.awssdk.services.s3.model.ListObjectsV2Response;
+import software.amazon.awssdk.services.s3.model.ObjectIdentifier;
 import software.amazon.awssdk.services.s3.model.Permission;
 import software.amazon.awssdk.services.s3.model.S3Exception;
+import software.amazon.awssdk.services.s3.model.S3Object;
 import software.amazon.awssdk.services.s3.model.Type;
 
 import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.HashMap;
+import java.util.List;
 import java.util.regex.Pattern;
+import java.util.stream.Collectors;
 
 /**
  *
@@ -96,16 +109,14 @@ public class S3DataSegmentPusherTest
     s3Client.upload(EasyMock.anyString(), EasyMock.anyString(), 
EasyMock.anyObject(File.class), EasyMock.anyObject(Grant.class));
     EasyMock.expectLastCall().once();
 
+    // nothing else is under the path, so there is nothing to clean up
+    
EasyMock.expect(s3Client.listObjectsV2(EasyMock.anyObject(ListObjectsV2Request.class)))
+            
.andReturn(ListObjectsV2Response.builder().isTruncated(false).build())
+            .once();
+
     EasyMock.replay(s3Client);
 
-    S3DataSegmentPusherConfig config = new S3DataSegmentPusherConfig()
-    {
-      @Override
-      public boolean isZip()
-      {
-        return false;
-      }
-    };
+    S3DataSegmentPusherConfig config = new S3DataSegmentPusherConfig();
     config.setBucket("bucket");
     config.setBaseKey("key");
     DataSegment segment = validate(
@@ -113,6 +124,7 @@ public class S3DataSegmentPusherTest
         "key/foo/2015-01-01T00:00:00\\.000Z_2016-01-01T00:00:00\\.000Z/0/0/",
         s3Client,
         config,
+        false,
         new byte[]{0x0, 0x0, 0x0, 0x1}
     );
     // V1 (test fixture) → not V10 → rangeable stamped as false (skips legacy 
HEAD probe).
@@ -133,16 +145,14 @@ public class S3DataSegmentPusherTest
     s3Client.upload(EasyMock.anyString(), EasyMock.anyString(), 
EasyMock.anyObject(File.class), EasyMock.anyObject(Grant.class));
     EasyMock.expectLastCall().once();
 
+    // nothing else is under the path, so there is nothing to clean up
+    
EasyMock.expect(s3Client.listObjectsV2(EasyMock.anyObject(ListObjectsV2Request.class)))
+            
.andReturn(ListObjectsV2Response.builder().isTruncated(false).build())
+            .once();
+
     EasyMock.replay(s3Client);
 
-    S3DataSegmentPusherConfig config = new S3DataSegmentPusherConfig()
-    {
-      @Override
-      public boolean isZip()
-      {
-        return false;
-      }
-    };
+    S3DataSegmentPusherConfig config = new S3DataSegmentPusherConfig();
     config.setBucket("bucket");
     config.setBaseKey("key");
 
@@ -152,12 +162,132 @@ public class S3DataSegmentPusherTest
         "key/foo/2015-01-01T00:00:00\\.000Z_2016-01-01T00:00:00\\.000Z/0/0/",
         s3Client,
         config,
+        false,
         new byte[]{0x0, 0x0, 0x0, 0x0A}
     );
     Assertions.assertEquals(10, (int) segment.getBinaryVersion());
     Assertions.assertEquals(Boolean.TRUE, 
segment.getLoadSpec().get("rangeable"));
   }
 
+  @Test
+  public void testPushNoZipRemovesObjectsLeftByPreviousPush() throws Exception
+  {
+    ServerSideEncryptingAmazonS3 s3Client = 
EasyMock.createStrictMock(ServerSideEncryptingAmazonS3.class);
+
+    Grant grant = Grant.builder()
+        
.grantee(Grantee.builder().id("ownerId").type(Type.CANONICAL_USER).build())
+        .permission(Permission.FULL_CONTROL)
+        .build();
+    
EasyMock.expect(s3Client.getBucketOwnerGrant(EasyMock.eq("bucket"))).andReturn(grant).once();
+
+    s3Client.upload(EasyMock.anyString(), EasyMock.anyString(), 
EasyMock.anyObject(File.class), EasyMock.anyObject(Grant.class));
+    EasyMock.expectLastCall().once();
+
+    // an index.zip from a zipped push and a smoosh chunk from a larger prior 
segment: nothing in the new segment
+    // references either, but the puller would list and download both
+    final String prefix = 
"key/foo/2015-01-01T00:00:00.000Z_2016-01-01T00:00:00.000Z/0/0/";
+    final String staleZip = prefix + "index.zip";
+    final String staleChunk = prefix + "00001.smoosh";
+    
EasyMock.expect(s3Client.listObjectsV2(EasyMock.anyObject(ListObjectsV2Request.class)))
+            .andReturn(
+                ListObjectsV2Response.builder()
+                                     .contents(
+                                         S3Object.builder().key(prefix + 
"version.bin").size(4L).build(),
+                                         
S3Object.builder().key(staleZip).size(128L).build(),
+                                         
S3Object.builder().key(staleChunk).size(64L).build()
+                                     )
+                                     .isTruncated(false)
+                                     .build()
+            )
+            .once();
+
+    final Capture<DeleteObjectsRequest> deleteRequest = Capture.newInstance();
+    EasyMock.expect(s3Client.deleteObjects(EasyMock.capture(deleteRequest)))
+            .andReturn(DeleteObjectsResponse.builder().build())
+            .once();
+
+    EasyMock.replay(s3Client);
+
+    S3DataSegmentPusherConfig config = new S3DataSegmentPusherConfig();
+    config.setBucket("bucket");
+    config.setBaseKey("key");
+    validate(
+        false,
+        "key/foo/2015-01-01T00:00:00\\.000Z_2016-01-01T00:00:00\\.000Z/0/0/",
+        s3Client,
+        config,
+        false,
+        new byte[]{0x0, 0x0, 0x0, 0x1}
+    );
+
+    // only the objects this push did not write are removed
+    Assertions.assertEquals(
+        ImmutableSet.of(staleZip, staleChunk),
+        
deleteRequest.getValue().delete().objects().stream().map(ObjectIdentifier::key).collect(Collectors.toSet())
+    );
+  }
+
+  @Test
+  public void testPushNoZipBatchesStaleObjectDeletesWithDefaultConfig() throws 
Exception
+  {
+    ServerSideEncryptingAmazonS3 s3Client = 
EasyMock.createStrictMock(ServerSideEncryptingAmazonS3.class);
+
+    Grant grant = Grant.builder()
+        
.grantee(Grantee.builder().id("ownerId").type(Type.CANONICAL_USER).build())
+        .permission(Permission.FULL_CONTROL)
+        .build();
+    
EasyMock.expect(s3Client.getBucketOwnerGrant(EasyMock.eq("bucket"))).andReturn(grant).once();
+
+    s3Client.upload(EasyMock.anyString(), EasyMock.anyString(), 
EasyMock.anyObject(File.class), EasyMock.anyObject(Grant.class));
+    EasyMock.expectLastCall().once();
+
+    // maxListingLength defaults to 1024, which is a legal listing size but 
more keys than DeleteObjects accepts, so
+    // the deletes have to be split even though the listing is not
+    final String prefix = 
"key/foo/2015-01-01T00:00:00.000Z_2016-01-01T00:00:00.000Z/0/0/";
+    final int staleCount = S3Utils.MAX_MULTI_OBJECT_DELETE_SIZE + 1;
+    final List<S3Object> listing = new ArrayList<>();
+    listing.add(S3Object.builder().key(prefix + 
"version.bin").size(4L).build());
+    for (int i = 0; i < staleCount; i++) {
+      listing.add(S3Object.builder().key(StringUtils.format("%s%05d.smoosh", 
prefix, i)).size(8L).build());
+    }
+    
EasyMock.expect(s3Client.listObjectsV2(EasyMock.anyObject(ListObjectsV2Request.class)))
+            
.andReturn(ListObjectsV2Response.builder().contents(listing).isTruncated(false).build())
+            .once();
+
+    final Capture<DeleteObjectsRequest> deleteRequests = 
Capture.newInstance(CaptureType.ALL);
+    EasyMock.expect(s3Client.deleteObjects(EasyMock.capture(deleteRequests)))
+            .andReturn(DeleteObjectsResponse.builder().build())
+            .times(2);
+
+    EasyMock.replay(s3Client);
+
+    S3DataSegmentPusherConfig config = new S3DataSegmentPusherConfig();
+    config.setBucket("bucket");
+    config.setBaseKey("key");
+    Assertions.assertEquals(1024, config.getMaxListingLength());
+
+    validate(
+        false,
+        "key/foo/2015-01-01T00:00:00\\.000Z_2016-01-01T00:00:00\\.000Z/0/0/",
+        s3Client,
+        config,
+        false,
+        new byte[]{0x0, 0x0, 0x0, 0x1}
+    );
+
+    // no request may carry more keys than S3 accepts, and between them they 
remove every stale object
+    for (DeleteObjectsRequest request : deleteRequests.getValues()) {
+      Assertions.assertTrue(
+          request.delete().objects().size() <= 
S3Utils.MAX_MULTI_OBJECT_DELETE_SIZE,
+          "DeleteObjects request carried " + request.delete().objects().size() 
+ " keys"
+      );
+    }
+    Assertions.assertEquals(
+        staleCount,
+        deleteRequests.getValues().stream().mapToInt(request -> 
request.delete().objects().size()).sum()
+    );
+  }
+
   @Test
   public void testPushZipDoesNotStampRangeable() throws Exception
   {
@@ -234,7 +364,7 @@ public class S3DataSegmentPusherTest
     config.setBucket("bucket");
     config.setBaseKey("key");
     // Default version.bin is V1 for historical reasons.
-    DataSegment segment = validate(useUniquePath, matcher, s3Client, config, 
new byte[]{0x0, 0x0, 0x0, 0x1});
+    DataSegment segment = validate(useUniquePath, matcher, s3Client, config, 
true, new byte[]{0x0, 0x0, 0x0, 0x1});
     Assertions.assertEquals(1, (int) segment.getBinaryVersion());
     return segment;
   }
@@ -244,10 +374,11 @@ public class S3DataSegmentPusherTest
       String matcher,
       ServerSideEncryptingAmazonS3 s3Client,
       S3DataSegmentPusherConfig config,
+      boolean zip,
       byte[] versionBytes
   ) throws IOException
   {
-    S3DataSegmentPusher pusher = new S3DataSegmentPusher(s3Client, config);
+    S3DataSegmentPusher pusher = new S3DataSegmentPusher(s3Client, config, new 
DeepStorageSegmentConfig(zip));
 
     // Create a mock segment on disk
     File tmp = new File(tempFolder, "version.bin");
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/CompactionTaskRunBase.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/CompactionTaskRunBase.java
index e6fe17fa852..34171324ba2 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/CompactionTaskRunBase.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/CompactionTaskRunBase.java
@@ -96,6 +96,7 @@ import org.apache.druid.segment.TestIndex;
 import org.apache.druid.segment.column.ColumnType;
 import org.apache.druid.segment.handoff.NoopSegmentHandoffNotifierFactory;
 import org.apache.druid.segment.join.NoopJoinableFactory;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import 
org.apache.druid.segment.loading.LeastBytesUsedStorageLocationSelectorStrategy;
 import org.apache.druid.segment.loading.LocalDataSegmentPuller;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
@@ -1731,7 +1732,7 @@ public abstract class CompactionTaskRunBase
         .config(new TaskConfigBuilder().build())
         .emitter(NoopServiceEmitter.instance())
         .taskActionClient(taskActionClient)
-        .segmentPusher(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig()))
+        .segmentPusher(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig(), new DeepStorageSegmentConfig()))
         .dataSegmentKiller(new NoopDataSegmentKiller())
         .joinableFactory(NoopJoinableFactory.INSTANCE)
         .segmentCacheManager(cacheManager)
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/IngestionTestBase.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/IngestionTestBase.java
index 17a40ed01d8..a63aca6bc25 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/IngestionTestBase.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/IngestionTestBase.java
@@ -72,6 +72,7 @@ import org.apache.druid.segment.SegmentSchemaMapping;
 import org.apache.druid.segment.TestIndex;
 import org.apache.druid.segment.incremental.RowIngestionMetersFactory;
 import org.apache.druid.segment.join.NoopJoinableFactory;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
 import org.apache.druid.segment.loading.SegmentCacheManager;
@@ -300,7 +301,7 @@ public abstract class IngestionTestBase extends 
InitializedNullHandlingTest
         .config(config)
         .taskExecutorNode(new DruidNode("druid/middlemanager", "localhost", 
false, 8091, null, true, false))
         .taskActionClient(createActionClient(task))
-        .segmentPusher(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig()))
+        .segmentPusher(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig(), new DeepStorageSegmentConfig()))
         .dataSegmentKiller(dataSegmentKiller)
         .joinableFactory(NoopJoinableFactory.INSTANCE)
         .jsonMapper(objectMapper)
@@ -482,7 +483,7 @@ public abstract class IngestionTestBase extends 
InitializedNullHandlingTest
             .config(config)
             .taskExecutorNode(new DruidNode("druid/middlemanager", 
"localhost", false, 8091, null, true, false))
             .taskActionClient(taskActionClient)
-            .segmentPusher(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig()))
+            .segmentPusher(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig(), new DeepStorageSegmentConfig()))
             .dataSegmentKiller(dataSegmentKiller)
             .joinableFactory(NoopJoinableFactory.INSTANCE)
             .jsonMapper(objectMapper)
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/AbstractParallelIndexSupervisorTaskTest.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/AbstractParallelIndexSupervisorTaskTest.java
index 6878235a242..d55e27a2990 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/AbstractParallelIndexSupervisorTaskTest.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/batch/parallel/AbstractParallelIndexSupervisorTaskTest.java
@@ -89,6 +89,7 @@ import 
org.apache.druid.segment.incremental.ParseExceptionReport;
 import org.apache.druid.segment.incremental.RowIngestionMetersFactory;
 import org.apache.druid.segment.incremental.RowIngestionMetersTotals;
 import org.apache.druid.segment.join.NoopJoinableFactory;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LocalDataSegmentPuller;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
@@ -684,7 +685,8 @@ public class AbstractParallelIndexSupervisorTaskTest 
extends IngestionTestBase
                   {
                     return localDeepStorage;
                   }
-                }
+                },
+                new DeepStorageSegmentConfig()
             )
         )
         .dataSegmentKiller(new NoopDataSegmentKiller())
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java
index f7fb80308b1..f0d08b0d941 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java
@@ -114,6 +114,7 @@ import 
org.apache.druid.segment.incremental.RowIngestionMetersTotals;
 import org.apache.druid.segment.indexing.DataSchema;
 import org.apache.druid.segment.join.NoopJoinableFactory;
 import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
 import org.apache.druid.segment.metadata.CentralizedDatasourceSchemaConfig;
@@ -664,8 +665,8 @@ public abstract class SeekableStreamIndexTaskTestBase 
extends EasyMockSupport
     };
     final LocalDataSegmentPusherConfig dataSegmentPusherConfig = new 
LocalDataSegmentPusherConfig();
     dataSegmentPusherConfig.storageDirectory = getSegmentDirectory();
-    dataSegmentPusherConfig.zip = true;
-    final DataSegmentPusher dataSegmentPusher = new 
LocalDataSegmentPusher(dataSegmentPusherConfig);
+    final DataSegmentPusher dataSegmentPusher =
+        new LocalDataSegmentPusher(dataSegmentPusherConfig, new 
DeepStorageSegmentConfig(true));
 
     toolboxFactory = new TaskToolboxFactory(
         null,
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/worker/shuffle/ShuffleDataSegmentPusherTest.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/worker/shuffle/ShuffleDataSegmentPusherTest.java
index bcd879ff8f5..3220c9c3972 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/worker/shuffle/ShuffleDataSegmentPusherTest.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/worker/shuffle/ShuffleDataSegmentPusherTest.java
@@ -38,6 +38,7 @@ import org.apache.druid.jackson.DefaultObjectMapper;
 import org.apache.druid.java.util.common.Intervals;
 import org.apache.druid.rpc.indexing.NoopOverlordClient;
 import org.apache.druid.rpc.indexing.OverlordClient;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LoadSpec;
 import org.apache.druid.segment.loading.LocalDataSegmentPuller;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
@@ -116,7 +117,9 @@ public class ShuffleDataSegmentPusherTest
                 {
                   return localDeepStore;
                 }
-              }));
+              },
+              new DeepStorageSegmentConfig()
+          ));
     }
     intermediaryDataManager.start();
     segmentPusher = new ShuffleDataSegmentPusher("supervisorTaskId", 
"subTaskId", intermediaryDataManager);
diff --git 
a/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQCompactionTaskRunTest.java
 
b/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQCompactionTaskRunTest.java
index 913bb35d072..5fafa602aea 100644
--- 
a/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQCompactionTaskRunTest.java
+++ 
b/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQCompactionTaskRunTest.java
@@ -95,6 +95,7 @@ import org.apache.druid.segment.indexing.TuningConfig;
 import org.apache.druid.segment.loading.AcquireSegmentAction;
 import org.apache.druid.segment.loading.AcquireSegmentResult;
 import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
 import org.apache.druid.segment.loading.SegmentCacheManager;
@@ -218,7 +219,7 @@ public class MSQCompactionTaskRunTest extends 
CompactionTaskRunBase
         binder -> 
binder.bind(PolicyEnforcer.class).toInstance(NoopPolicyEnforcer.instance()),
         binder -> binder.bind(WireTransferableContext.class).toInstance(new 
WireTransferableContext(null, null, true)),
         binder -> binder.bind(DataSegmentPusher.class)
-                        .toInstance(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig())),
+                        .toInstance(new LocalDataSegmentPusher(new 
LocalDataSegmentPusherConfig(), new DeepStorageSegmentConfig())),
         binder -> 
binder.bind(DataServerQueryHandlerFactory.class).toProvider(Providers.of(null)),
         binder -> binder.bind(Escalator.class).toProvider(Providers.of(null)),
         binder -> binder.bind(QueryProcessingPool.class)
diff --git 
a/multi-stage-query/src/test/java/org/apache/druid/msq/test/CalciteMSQTestsHelper.java
 
b/multi-stage-query/src/test/java/org/apache/druid/msq/test/CalciteMSQTestsHelper.java
index 5827a71c50f..16a1c0bd348 100644
--- 
a/multi-stage-query/src/test/java/org/apache/druid/msq/test/CalciteMSQTestsHelper.java
+++ 
b/multi-stage-query/src/test/java/org/apache/druid/msq/test/CalciteMSQTestsHelper.java
@@ -50,6 +50,7 @@ import org.apache.druid.query.QueryProcessingPool;
 import org.apache.druid.query.groupby.TestGroupByBuffers;
 import org.apache.druid.segment.TestIndex;
 import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
 import org.apache.druid.segment.loading.SegmentCacheManager;
@@ -121,7 +122,10 @@ public class CalciteMSQTestsHelper
         TestSegmentManager testSegmentManager
     )
     {
-      return new MSQTestDelegateDataSegmentPusher(new 
LocalDataSegmentPusher(config), testSegmentManager);
+      return new MSQTestDelegateDataSegmentPusher(
+          new LocalDataSegmentPusher(config, new DeepStorageSegmentConfig()),
+          testSegmentManager
+      );
     }
 
     @Provides
diff --git 
a/multi-stage-query/src/test/java/org/apache/druid/msq/test/MSQTestBase.java 
b/multi-stage-query/src/test/java/org/apache/druid/msq/test/MSQTestBase.java
index 0f6f80eeeb6..45b1eda4ce4 100644
--- a/multi-stage-query/src/test/java/org/apache/druid/msq/test/MSQTestBase.java
+++ b/multi-stage-query/src/test/java/org/apache/druid/msq/test/MSQTestBase.java
@@ -160,6 +160,7 @@ import org.apache.druid.segment.column.RowSignature;
 import org.apache.druid.segment.incremental.IncrementalIndexSchema;
 import org.apache.druid.segment.loading.AcquireMode;
 import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import 
org.apache.druid.segment.loading.LeastBytesUsedStorageLocationSelectorStrategy;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
@@ -534,7 +535,7 @@ public class MSQTestBase extends BaseCalciteQueryTest
           LocalDataSegmentPusherConfig config = new 
LocalDataSegmentPusherConfig();
           config.storageDirectory = newTempFolder("storageDir");
           binder.bind(DataSegmentPusher.class).toInstance(new 
MSQTestDelegateDataSegmentPusher(
-              new LocalDataSegmentPusher(config),
+              new LocalDataSegmentPusher(config, new 
DeepStorageSegmentConfig()),
               testSegmentManager
           ));
           binder.bind(DataSegmentAnnouncer.class).toInstance(new 
NoopDataSegmentAnnouncer());
diff --git 
a/processing/src/main/java/org/apache/druid/segment/loading/DeepStorageSegmentConfig.java
 
b/processing/src/main/java/org/apache/druid/segment/loading/DeepStorageSegmentConfig.java
new file mode 100644
index 00000000000..9ae3341d6db
--- /dev/null
+++ 
b/processing/src/main/java/org/apache/druid/segment/loading/DeepStorageSegmentConfig.java
@@ -0,0 +1,73 @@
+/*
+ * 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.segment.loading;
+
+import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonProperty;
+
+import javax.annotation.Nullable;
+
+/**
+ * The {@code druid.storage} properties that describe how segments are laid 
out in deep storage rather than how to
+ * reach it, and so are shared by every {@link DataSegmentPusher} instead of 
belonging to any one of them. Currently
+ * that is only {@code druid.storage.zip}, but anything else that is a 
property of the layout rather than of a
+ * particular deep storage belongs here too.
+ * <p>
+ * This is bound once in core, by {@code LocalDataStorageDruidModule}, instead 
of once per deep storage
+ * implementation, so that these properties mean the same thing whatever 
{@code druid.storage.type} is set to, and so
+ * that each implementation does not have to redeclare them in its own config 
class. Any single-implementation
+ * property still belongs in that implementation's own config, next to its 
bucket and credentials.
+ * <p>
+ * Values here are deliberately tri-state: when a property is unset, the 
getter falls back to the default the calling
+ * implementation passes in, since those defaults differ per implementation.
+ */
+public class DeepStorageSegmentConfig
+{
+  @JsonProperty
+  @Nullable
+  private final Boolean zip;
+
+  @JsonCreator
+  public DeepStorageSegmentConfig(@JsonProperty("zip") @Nullable Boolean zip)
+  {
+    this.zip = zip;
+  }
+
+  public DeepStorageSegmentConfig()
+  {
+    this(null);
+  }
+
+  /**
+   * Whether segments should be written as a zip file, falling back to {@code 
defaultZip} when
+   * {@code druid.storage.zip} has not been set.
+   */
+  public boolean isZip(boolean defaultZip)
+  {
+    return zip == null ? defaultZip : zip;
+  }
+
+  @JsonProperty("zip")
+  @Nullable
+  public Boolean getZip()
+  {
+    return zip;
+  }
+}
diff --git 
a/processing/src/test/java/org/apache/druid/segment/loading/DeepStorageSegmentConfigTest.java
 
b/processing/src/test/java/org/apache/druid/segment/loading/DeepStorageSegmentConfigTest.java
new file mode 100644
index 00000000000..89e65ea3228
--- /dev/null
+++ 
b/processing/src/test/java/org/apache/druid/segment/loading/DeepStorageSegmentConfigTest.java
@@ -0,0 +1,105 @@
+/*
+ * 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.segment.loading;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import jakarta.validation.Validation;
+import jakarta.validation.Validator;
+import org.apache.druid.guice.JsonConfigurator;
+import org.apache.druid.jackson.DefaultObjectMapper;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.util.Properties;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class DeepStorageSegmentConfigTest
+{
+  private static final String PROP_PREFIX = "druid.storage";
+
+  private final ObjectMapper jsonMapper = new DefaultObjectMapper();
+  private final Validator validator = 
Validation.buildDefaultValidatorFactory().getValidator();
+
+  @Test
+  public void testUnsetFallsBackToCallerDefault()
+  {
+    final DeepStorageSegmentConfig config = new DeepStorageSegmentConfig();
+    assertNull(config.getZip());
+    // each DataSegmentPusher keeps its own default when druid.storage.zip is 
not set
+    assertTrue(config.isZip(true));
+    assertFalse(config.isZip(false));
+  }
+
+  @Test
+  public void testSetOverridesCallerDefault()
+  {
+    assertTrue(new DeepStorageSegmentConfig(true).isZip(false));
+    assertFalse(new DeepStorageSegmentConfig(false).isZip(true));
+  }
+
+  @Test
+  public void testConfigurateFromProperties()
+  {
+    final Properties properties = new Properties();
+    properties.setProperty("druid.storage.zip", "false");
+
+    final DeepStorageSegmentConfig config = configurate(properties);
+    assertEquals(Boolean.FALSE, config.getZip());
+    assertFalse(config.isZip(true));
+  }
+
+  @Test
+  public void testConfigurateIgnoresOtherStorageProperties()
+  {
+    // druid.storage is shared with the deep storage implementation's own 
config, so the properties belonging to it
+    // must not upset this one
+    final Properties properties = new Properties();
+    properties.setProperty("druid.storage.type", "s3");
+    properties.setProperty("druid.storage.bucket", "bucket");
+    properties.setProperty("druid.storage.baseKey", "baseKey");
+
+    final DeepStorageSegmentConfig config = configurate(properties);
+    assertNull(config.getZip());
+    assertTrue(config.isZip(true));
+  }
+
+  @Test
+  public void testSerde() throws IOException
+  {
+    assertEquals(
+        Boolean.TRUE,
+        jsonMapper.readValue(jsonMapper.writeValueAsBytes(new 
DeepStorageSegmentConfig(true)), DeepStorageSegmentConfig.class)
+                  .getZip()
+    );
+    assertNull(
+        jsonMapper.readValue(jsonMapper.writeValueAsBytes(new 
DeepStorageSegmentConfig()), DeepStorageSegmentConfig.class)
+                  .getZip()
+    );
+  }
+
+  private DeepStorageSegmentConfig configurate(Properties properties)
+  {
+    return new JsonConfigurator(jsonMapper, validator).configurate(properties, 
PROP_PREFIX, DeepStorageSegmentConfig.class);
+  }
+}
diff --git 
a/server/src/main/java/org/apache/druid/guice/LocalDataStorageDruidModule.java 
b/server/src/main/java/org/apache/druid/guice/LocalDataStorageDruidModule.java
index 8f1d44c3fe3..7109a08fb36 100644
--- 
a/server/src/main/java/org/apache/druid/guice/LocalDataStorageDruidModule.java
+++ 
b/server/src/main/java/org/apache/druid/guice/LocalDataStorageDruidModule.java
@@ -29,6 +29,7 @@ import org.apache.druid.data.SearchableVersionedDataFinder;
 import org.apache.druid.initialization.DruidModule;
 import org.apache.druid.segment.loading.DataSegmentKiller;
 import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
 import org.apache.druid.segment.loading.LocalDataSegmentKiller;
 import org.apache.druid.segment.loading.LocalDataSegmentPusher;
 import org.apache.druid.segment.loading.LocalDataSegmentPusherConfig;
@@ -48,6 +49,11 @@ public class LocalDataStorageDruidModule implements 
DruidModule
   {
     bindDeepStorageLocal(binder);
 
+    // The segment layout properties are not specific to local deep storage: 
they are bound here, next to the
+    // DataSegmentPusher choice, because this module is always installed, so 
every deep storage implementation can read
+    // the same config object instead of declaring the properties itself.
+    JsonConfigProvider.bind(binder, "druid.storage", 
DeepStorageSegmentConfig.class);
+
     PolyBind.createChoice(
         binder,
         "druid.storage.type",
diff --git 
a/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusher.java
 
b/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusher.java
index 584fa184a7c..496c134b464 100644
--- 
a/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusher.java
+++ 
b/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusher.java
@@ -44,12 +44,19 @@ public class LocalDataSegmentPusher implements 
DataSegmentPusher
   public static final String INDEX_DIR = "index";
   public static final String INDEX_ZIP_FILENAME = "index.zip";
 
+  /**
+   * Local deep storage writes segments as a directory of files unless {@code 
druid.storage.zip} says otherwise.
+   */
+  private static final boolean DEFAULT_ZIP = false;
+
   private final LocalDataSegmentPusherConfig config;
+  private final boolean zip;
 
   @Inject
-  public LocalDataSegmentPusher(LocalDataSegmentPusherConfig config)
+  public LocalDataSegmentPusher(LocalDataSegmentPusherConfig config, 
DeepStorageSegmentConfig deepStorageConfig)
   {
     this.config = config;
+    this.zip = deepStorageConfig.isZip(DEFAULT_ZIP);
   }
 
   @Override
@@ -80,7 +87,7 @@ public class LocalDataSegmentPusher implements 
DataSegmentPusher
       return segment.withLoadSpec(makeLoadSpec(outDir.toURI()))
                     .withSize(size)
                     
.withBinaryVersion(SegmentUtils.getVersionFromDir(dataSegmentFile));
-    } else if (config.isZip()) {
+    } else if (zip) {
       return pushZip(dataSegmentFile, outDir, segment);
     } else {
       return pushNoZip(dataSegmentFile, outDir, segment);
diff --git 
a/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusherConfig.java
 
b/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusherConfig.java
index e539ad20af6..15868da4de1 100644
--- 
a/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusherConfig.java
+++ 
b/server/src/main/java/org/apache/druid/segment/loading/LocalDataSegmentPusherConfig.java
@@ -30,16 +30,8 @@ public class LocalDataSegmentPusherConfig
   @JsonProperty
   public File storageDirectory = new File("/tmp/druid/localStorage");
 
-  @JsonProperty
-  public boolean zip = false;
-
   public File getStorageDirectory()
   {
     return storageDirectory;
   }
-
-  public boolean isZip()
-  {
-    return zip;
-  }
 }
diff --git 
a/server/src/test/java/org/apache/druid/guice/LocalDataStorageDruidModuleTest.java
 
b/server/src/test/java/org/apache/druid/guice/LocalDataStorageDruidModuleTest.java
index ce69a34b4ac..41edf56dea2 100644
--- 
a/server/src/test/java/org/apache/druid/guice/LocalDataStorageDruidModuleTest.java
+++ 
b/server/src/test/java/org/apache/druid/guice/LocalDataStorageDruidModuleTest.java
@@ -20,18 +20,30 @@
 package org.apache.druid.guice;
 
 import com.google.common.collect.ImmutableList;
+import com.google.common.io.Files;
+import com.google.common.primitives.Ints;
 import com.google.inject.Injector;
 import com.google.inject.Module;
 import com.google.inject.TypeLiteral;
+import org.apache.druid.java.util.common.FileUtils;
+import org.apache.druid.java.util.common.Intervals;
 import org.apache.druid.segment.column.ColumnConfig;
+import org.apache.druid.segment.loading.DataSegmentPusher;
 import org.apache.druid.segment.loading.OmniDataSegmentKiller;
 import org.apache.druid.segment.loading.RandomStorageLocationSelectorStrategy;
 import org.apache.druid.segment.loading.StorageLocation;
 import org.apache.druid.segment.loading.StorageLocationSelectorStrategy;
+import org.apache.druid.timeline.DataSegment;
+import org.apache.druid.timeline.SegmentId;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
 
+import javax.annotation.Nullable;
+import java.io.File;
+import java.io.IOException;
 import java.util.List;
+import java.util.Properties;
 
 public class LocalDataStorageDruidModuleTest
 {
@@ -47,6 +59,49 @@ public class LocalDataStorageDruidModuleTest
     );
   }
 
+  /**
+   * {@code druid.storage.zip} is bound by this module for every deep storage 
implementation to read, so check that it
+   * actually reaches the pusher the module provisions, rather than only that 
it parses.
+   */
+  @Test
+  public void testDataSegmentPusherHonorsStorageZip(@TempDir File tempDir) 
throws IOException
+  {
+    Assertions.assertTrue(new File(pushSegment(tempDir, "true"), 
"index.zip").isFile());
+  }
+
+  @Test
+  public void testDataSegmentPusherDefaultsToUnzipped(@TempDir File tempDir) 
throws IOException
+  {
+    // local deep storage writes a directory of files when the property is 
unset, which is the default the tri-state
+    // DeepStorageSegmentConfig exists to preserve
+    Assertions.assertTrue(new File(pushSegment(tempDir, null), 
"index").isDirectory());
+  }
+
+  /**
+   * Pushes a one-file segment through a {@link DataSegmentPusher} provisioned 
by this module, and returns the
+   * directory it was pushed to.
+   */
+  private static File pushSegment(File tempDir, @Nullable String zip) throws 
IOException
+  {
+    final File storageDir = new File(tempDir, "deepStorage");
+    final File segmentDir = new File(tempDir, "segment");
+    FileUtils.mkdirp(segmentDir);
+    Files.asByteSink(new File(segmentDir, 
"version.bin")).write(Ints.toByteArray(0x9));
+
+    final Injector injector = createInjector();
+    // JsonConfigProvider reads these lazily, on the first getInstance below
+    final Properties properties = injector.getInstance(Properties.class);
+    properties.setProperty("druid.storage.storageDirectory", 
storageDir.getAbsolutePath());
+    if (zip != null) {
+      properties.setProperty("druid.storage.zip", zip);
+    }
+
+    final DataSegmentPusher pusher = 
injector.getInstance(DataSegmentPusher.class);
+    final DataSegment segment = DataSegment.builder(SegmentId.of("ds", 
Intervals.utc(0, 1), "v1", 0)).build();
+
+    return new File(storageDir, pusher.getStorageDir(pusher.push(segmentDir, 
segment, false), false));
+  }
+
   private static Injector createInjector()
   {
     return GuiceInjectors.makeStartupInjectorWithModules(
diff --git 
a/server/src/test/java/org/apache/druid/segment/loading/LocalDataSegmentPusherTest.java
 
b/server/src/test/java/org/apache/druid/segment/loading/LocalDataSegmentPusherTest.java
index 1fa828b8183..d090161e572 100644
--- 
a/server/src/test/java/org/apache/druid/segment/loading/LocalDataSegmentPusherTest.java
+++ 
b/server/src/test/java/org/apache/druid/segment/loading/LocalDataSegmentPusherTest.java
@@ -75,14 +75,12 @@ public class LocalDataSegmentPusherTest
   public void setUp() throws IOException
   {
     config = new LocalDataSegmentPusherConfig();
-    config.zip = false;
     config.storageDirectory = temporaryFolder.newFolder();
-    localDataSegmentPusher = new LocalDataSegmentPusher(config);
+    localDataSegmentPusher = new LocalDataSegmentPusher(config, new 
DeepStorageSegmentConfig(false));
 
     configZip = new LocalDataSegmentPusherConfig();
-    configZip.zip = true;
     configZip.storageDirectory = temporaryFolder.newFolder();
-    localDataSegmentPusherZip = new LocalDataSegmentPusher(configZip);
+    localDataSegmentPusherZip = new LocalDataSegmentPusher(configZip, new 
DeepStorageSegmentConfig(true));
 
     dataSegmentFiles = temporaryFolder.newFolder();
     Files.asByteSink(new File(dataSegmentFiles, 
"version.bin")).write(Ints.toByteArray(0x9));


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

Reply via email to