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

kfaraz pushed a commit to branch 38.0.0
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/38.0.0 by this push:
     new 9af3e3f47e2 fix: Stop kill task if segment IDs that share load spec of 
killable segment could not be determined (#19737)
9af3e3f47e2 is described below

commit 9af3e3f47e29eed23945004425d477d9d040706d
Author: Kashif Faraz <[email protected]>
AuthorDate: Sat Jul 25 19:21:16 2026 +0530

    fix: Stop kill task if segment IDs that share load spec of killable segment 
could not be determined (#19737)
    
    Bug
    ---
    When concurrent append and replace is enabled and `kill` tasks (or embedded 
kill on Overlord)
    are used, there may be potential data loss if the task action 
`retrieveUpgradedFromSegmentIds`
    or `retrieveUpgradedToSegmentIds` fired by the `kill` task fails. This is 
because the failure of
    these task actions is currently ignored and we may end up removing segment 
files from deep store
    even if the same load specs are shared by other segments.
    
    Changes
    -------
    - Stop `kill` task (and embedded kill tasks) if parent IDs of the killable 
unused segments could not
    be determined. There are 2 cases possible:
      - Segment does not share load spec: Do not remove segment from metadata 
store, otherwise deep
    store files would become orphaned and will never be removed.
      - Segment shares load spec with other segments: Do not remove segment 
files from deep store as
    the same files are needed by other segments.
    - Stop `kill` task (and embedded kill tasks) if sibling segment IDs or 
children segment IDs could not
    be determined. Same cases as above.
    - Improve performance of embedded kill tasks by avoiding extra DB call to 
fetch parent IDs
      - pre-fetch the parent IDs in the previous call which identifies the 
killable unused segments in the first place.
    
    (cherry picked from commit ff7584596d8d4aec6a0fdb21bbb9ca6ba96c7f39)
---
 docs/data-management/delete.md                     |   8 +-
 .../common/task/KillUnusedSegmentsTask.java        | 162 +++++++++++++--------
 .../overlord/duty/UnusedSegmentsKiller.java        |  25 +++-
 .../indexing/common/actions/TaskActionTestKit.java |  36 +++++
 .../indexing/common/task/IngestionTestBase.java    |   2 +-
 .../common/task/KillUnusedSegmentsTaskTest.java    | 150 ++++++++++++++++---
 .../druid/indexing/overlord/TaskLifecycleTest.java |   2 +-
 .../overlord/duty/UnusedSegmentsKillerTest.java    |  32 +++-
 .../TestIndexerMetadataStorageCoordinator.java     |   2 +-
 .../IndexerMetadataStorageCoordinator.java         |   5 +-
 .../IndexerSQLMetadataStorageCoordinator.java      |   2 +-
 .../druid/metadata/SqlSegmentsMetadataQuery.java   |  26 +++-
 .../IndexerSQLMetadataStorageCoordinatorTest.java  |   9 +-
 13 files changed, 364 insertions(+), 97 deletions(-)

diff --git a/docs/data-management/delete.md b/docs/data-management/delete.md
index cf571566ede..799e8b4b8b9 100644
--- a/docs/data-management/delete.md
+++ b/docs/data-management/delete.md
@@ -112,9 +112,13 @@ Some of the parameters used in the task payload are 
further explained below:
 | `limit`     | null (no limit) | Maximum number of segments for the kill task 
to delete.|
 | `maxUsedStatusLastUpdatedTime` | null (no cutoff) | Maximum timestamp used 
as a cutoff to include unused segments. The kill task only considers segments 
which lie in the specified `interval` and were marked as unused no later than 
this time. The default behavior is to kill all unused segments in the 
`interval` regardless of when they where marked as unused.|
 
-
-**WARNING:** The `kill` task permanently removes all information about the 
affected segments from the metadata store and
+:::warning
+- The `kill` task permanently removes all information about the affected 
segments from the metadata store and
 deep storage. This operation cannot be undone.
+- When using [concurrent locks](../ingestion/concurrent-append-replace.md) to 
run a `kill` task, ensure to keep a large
+enough buffer period before killing segments after they have been marked as 
unused. Otherwise, there may be a potential
+data loss if a concurrent append job upgrades one of the segments that are 
being killed.
+:::
 
 ### Auto-kill data using Coordinator duties
 
diff --git 
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java
 
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java
index fe58c264ca8..21b9e3f6a79 100644
--- 
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java
+++ 
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java
@@ -48,10 +48,10 @@ import org.apache.druid.java.util.common.ISE;
 import org.apache.druid.java.util.common.StringUtils;
 import org.apache.druid.java.util.common.logger.Logger;
 import org.apache.druid.server.coordination.BroadcastDatasourceLoadingSpec;
+import org.apache.druid.server.http.DataSegmentPlus;
 import org.apache.druid.server.lookup.cache.LookupLoadingSpec;
 import org.apache.druid.server.security.ResourceAction;
 import org.apache.druid.timeline.DataSegment;
-import org.apache.druid.timeline.SegmentId;
 import org.apache.druid.utils.CollectionUtils;
 import org.joda.time.DateTime;
 import org.joda.time.Interval;
@@ -67,10 +67,10 @@ import java.util.Map;
 import java.util.NavigableMap;
 import java.util.Set;
 import java.util.TreeMap;
+import java.util.function.Function;
 import java.util.stream.Collectors;
 
 /**
- * <p/>
  * The client representation of this task is {@link 
ClientKillUnusedSegmentsTaskQuery}.
  * JSON serialization fields of this class must correspond to those of {@link
  * ClientKillUnusedSegmentsTaskQuery}, except for {@link #id} and {@link 
#context} fields.
@@ -87,6 +87,10 @@ import java.util.stream.Collectors;
  * <li> Filter the set of unreferenced segments using load specs from the set 
of used segments. </li>
  * <li> Kill the filtered set of segments from deep storage. </li>
  * </ol>
+ * Note: When {@link Tasks#USE_CONCURRENT_LOCKS} is true, keep a large buffer
+ * period before killing segments after they have been marked as unused.
+ * Otherwise, there may be a potential data loss if a concurrent APPEND job
+ * upgrades one of the segments that are being killed.
  */
 public class KillUnusedSegmentsTask extends AbstractFixedIntervalTask
 {
@@ -211,7 +215,7 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
     // List unused segments
     int nextBatchSize = computeNextBatchSize(numSegmentsKilled);
     @Nullable Integer numTotalBatches = getNumTotalBatches();
-    List<DataSegment> unusedSegments;
+    List<DataSegmentPlus> unusedSegmentsPlus;
     logInfo(
         "Starting kill for datasource[%s] in interval[%s] and versions[%s] 
with batchSize[%d], up to limit[%d]"
         + " segments before maxUsedStatusLastUpdatedTime[%s] will be 
deleted%s",
@@ -236,12 +240,25 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
         break;
       }
 
-      unusedSegments = fetchNextBatchOfUnusedSegments(toolbox, nextBatchSize);
+      unusedSegmentsPlus = fetchNextBatchOfUnusedSegments(toolbox, 
nextBatchSize);
+      if (unusedSegmentsPlus.isEmpty()) {
+        // No more segments eligible for kill, do not proceed further
+        break;
+      }
 
       // Fetch locks each time as a revokal could have occurred in between 
batches
       final NavigableMap<DateTime, List<TaskLock>> taskLockMap
               = getNonRevokedTaskLockMap(toolbox.getTaskActionClient());
 
+      final Set<DataSegment> unusedSegments = unusedSegmentsPlus.stream()
+                                                                
.map(DataSegmentPlus::getDataSegment)
+                                                                
.collect(Collectors.toSet());
+      final Map<String, DataSegmentPlus> unusedIdToSegmentPlus = 
CollectionUtils.toMap(
+          unusedSegmentsPlus,
+          segment -> segment.getDataSegment().getId().toString(),
+          Function.identity()
+      );
+
       if (!TaskLocks.isLockCoversSegments(taskLockMap, unusedSegments)) {
         throw new ISE(
                 "Locks[%s] for task[%s] can't cover segments[%s]",
@@ -251,62 +268,46 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
         );
       }
 
-      // Kill segments. Order is important here:
-      // Retrieve the segment upgrade infos for the batch _before_ the 
segments are nuked
-      // We then want the nuke action to clean up the metadata records 
_before_ the segments are removed from storage.
-      // This helps maintain that we will always have a storage segment if the 
metadata segment is present.
-      // Determine the subset of segments to be killed from deep storage based 
on loadspecs.
-      // If the segment nuke throws an exception, then the segment cleanup is 
abandoned.
-
-      // Determine upgraded segment ids before nuking
-      final Set<String> segmentIds = unusedSegments.stream()
-                                                   .map(DataSegment::getId)
-                                                   .map(SegmentId::toString)
-                                                   
.collect(Collectors.toSet());
-      final Map<String, String> upgradedFromSegmentIds = new HashMap<>();
-      try {
-        upgradedFromSegmentIds.putAll(
-            taskActionClient.submit(
-                new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), 
segmentIds)
-            ).getUpgradedFromSegmentIds()
-        );
-      }
-      catch (Exception e) {
-        LOG.warn(
-            e,
-            "Could not retrieve parent segment ids using task 
action[retrieveUpgradedFromSegmentIds]."
-            + " Overlord may be on an older version."
-        );
-      }
+      // Kill segments - order of steps 1, 2, 3, 4 must remain the same
 
-      // Nuke Segments
-      taskActionClient.submit(new SegmentNukeAction(new 
HashSet<>(unusedSegments)));
-      emitMetric(toolbox.getEmitter(), 
TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size());
+      // 1. Determine parent segment ids of killable unused segments
+      final Map<String, String> upgradedFromSegmentIds
+          = fetchParentIdsForSegments(toolbox, unusedIdToSegmentPlus);
 
-      // Determine segments to be killed
-      final List<DataSegment> segmentsToBeKilled
-          = getKillableSegments(unusedSegments, upgradedFromSegmentIds, 
usedSegmentLoadSpecs, taskActionClient);
+      // 2. Identify killable segments whose load specs are not shared with 
any other segment
+      final List<DataSegment> segmentsToKillFromDeepStore = 
getKillableSegments(
+          unusedIdToSegmentPlus,
+          upgradedFromSegmentIds,
+          usedSegmentLoadSpecs,
+          taskActionClient
+      );
 
+      // 2a. Track segments that cannot be removed from deep store yet
       final Set<DataSegment> segmentsNotKilled = new HashSet<>(unusedSegments);
-      segmentsToBeKilled.forEach(segmentsNotKilled::remove);
-
+      segmentsToKillFromDeepStore.forEach(segmentsNotKilled::remove);
       if (!segmentsNotKilled.isEmpty()) {
         LOG.warn(
-            "Skipping kill of [%d] segments from deep storage as their load 
specs are used by other segments.",
-            segmentsNotKilled.size()
+            "Skipping kill of [%d] segments of datasource[%s] from deep 
storage"
+            + " as their load specs are shared by other segments.",
+            segmentsNotKilled.size(), getDataSource()
         );
       }
 
-      toolbox.getDataSegmentKiller().kill(segmentsToBeKilled);
-      emitMetric(toolbox.getEmitter(), 
TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE, segmentsToBeKilled.size());
+      // 3. Nuke all eligible unused segments
+      taskActionClient.submit(new SegmentNukeAction(unusedSegments));
+      emitMetric(toolbox.getEmitter(), 
TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedIdToSegmentPlus.size());
+
+      // 4. Delete deep store files only for segments which do not share load 
specs with other segments
+      toolbox.getDataSegmentKiller().kill(segmentsToKillFromDeepStore);
+      emitMetric(toolbox.getEmitter(), 
TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE, 
segmentsToKillFromDeepStore.size());
 
       numBatchesProcessed++;
-      numSegmentsKilled += segmentsToBeKilled.size();
+      numSegmentsKilled += segmentsToKillFromDeepStore.size();
 
       logInfo("Processed [%d] batches for kill task[%s].", 
numBatchesProcessed, getId());
 
       nextBatchSize = computeNextBatchSize(numSegmentsKilled);
-    } while (!unusedSegments.isEmpty() && (null == numTotalBatches || 
numBatchesProcessed < numTotalBatches));
+    } while (!unusedSegmentsPlus.isEmpty() && (null == numTotalBatches || 
numBatchesProcessed < numTotalBatches));
 
     final String taskId = getId();
     logInfo(
@@ -342,7 +343,7 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
   /**
    * Fetches the next batch of unused segments that are eligible for kill.
    */
-  protected List<DataSegment> fetchNextBatchOfUnusedSegments(TaskToolbox 
toolbox, int nextBatchSize) throws IOException
+  protected List<DataSegmentPlus> fetchNextBatchOfUnusedSegments(TaskToolbox 
toolbox, int nextBatchSize) throws IOException
   {
     return toolbox.getTaskActionClient().submit(
         new RetrieveUnusedSegmentsAction(
@@ -352,7 +353,44 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
             nextBatchSize,
             maxUsedStatusLastUpdatedTime
         )
-    );
+    )
+                  .stream()
+                  .map(segment -> new DataSegmentPlus(segment, null, null, 
null, null, null, null, null))
+                  .collect(Collectors.toList());
+  }
+
+  /**
+   * Fetches the parent IDs (if any) for the given unused segments.
+   *
+   * @param unusedIdToSegmentPlus Map containing unused segments whose parent 
IDs
+   *                              need to be fetched
+   * @return Map from segment ID to the segment ID from which
+   * it was upgraded. If an input segment was not upgraded from any other 
segment,
+   * it does not have an entry in the map.
+   */
+  protected Map<String, String> fetchParentIdsForSegments(
+      TaskToolbox toolbox,
+      Map<String, DataSegmentPlus> unusedIdToSegmentPlus
+  )
+  {
+    try {
+      return toolbox.getTaskActionClient().submit(
+          new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), 
unusedIdToSegmentPlus.keySet())
+      ).getUpgradedFromSegmentIds();
+    }
+    catch (Exception e) {
+      // Do not proceed with killing these segments as we cannot be sure if 
their
+      // load spec is shared by any other segment or not. If load spec is 
shared,
+      // segment files cannot be deleted from deep store. If load spec is not
+      // shared, segments cannot be deleted from metadata store as that would
+      // leave deep store files orphaned, and they would never be cleaned up.
+      throw new ISE(
+          e,
+          "Could not retrieve parent segment ids using task 
action[retrieveUpgradedFromSegmentIds]."
+          + " Stopping kill task to avoid data loss in case the segment files"
+          + " are shared by other segments."
+      );
+    }
   }
 
   /**
@@ -387,21 +425,20 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
    * @return list of segments to kill from deep storage
    */
   private List<DataSegment> getKillableSegments(
-      List<DataSegment> unusedSegments,
+      Map<String, DataSegmentPlus> unusedSegments,
       Map<String, String> upgradedFromSegmentIds,
       Set<Map<String, Object>> usedSegmentLoadSpecs,
       TaskActionClient taskActionClient
   )
   {
-
-    // Determine parentId for each unused segment
+    // Determine parentId (or self, if no parent) for each unused segment
     final Map<String, Set<DataSegment>> parentIdToUnusedSegments = new 
HashMap<>();
-    for (DataSegment segment : unusedSegments) {
-      final String segmentId = segment.getId().toString();
+    for (Map.Entry<String, DataSegmentPlus> entry : unusedSegments.entrySet()) 
{
+      final String segmentId = entry.getKey();
       parentIdToUnusedSegments.computeIfAbsent(
           upgradedFromSegmentIds.getOrDefault(segmentId, segmentId),
           k -> new HashSet<>()
-      ).add(segment);
+      ).add(entry.getValue().getDataSegment());
     }
 
     // Check if the parent or any of its children exist in metadata store
@@ -411,10 +448,11 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
       );
       if (response != null && response.getUpgradedToSegmentIds() != null) {
         response.getUpgradedToSegmentIds().forEach((parent, children) -> {
-          if (!CollectionUtils.isNullOrEmpty(children)) {
-            // Do not kill segment if its parent or any of its siblings still 
exist in metadata store
+          if (!unusedSegments.keySet().containsAll(children)) {
+            // Do not kill segment if its load spec is shared by another 
segment
+            // which is not being killed.
             LOG.info(
-                "Skipping kill of segments[%s] as its load spec is also used 
by segment IDs[%s].",
+                "Skipping kill of segments[%s] as its load spec is shared by 
segment IDs[%s].",
                 parentIdToUnusedSegments.get(parent), children
             );
             parentIdToUnusedSegments.remove(parent);
@@ -423,10 +461,14 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
       }
     }
     catch (Exception e) {
-      LOG.warn(
+      // Do not proceed with the kill of any segment as we cannot be sure if 
their
+      // load specs are shared by any other segment
+      throw new ISE(
           e,
-          "Could not retrieve referenced ids using task 
action[retrieveUpgradedToSegmentIds]."
-          + " Overlord may be on an older version."
+          "Could not perform task action[retrieveUpgradedToSegmentIds] to 
retrieve"
+          + " segment IDs which share load specs with segments being killed."
+          + " Stopping kill task to avoid data loss in case the segment files"
+          + " are shared by other segments."
       );
     }
 
@@ -449,7 +491,7 @@ public class KillUnusedSegmentsTask extends 
AbstractFixedIntervalTask
   {
     boolean isPresent = usedSegmentLoadSpecs.contains(segment.getLoadSpec());
     if (isPresent) {
-      LOG.info("Skipping kill of segment[%s] as its load spec is also used by 
other segments.", segment);
+      LOG.info("Skipping kill of segment[%s] as its load spec is shared by 
other 'used' segments.", segment);
     }
     return isPresent;
   }
diff --git 
a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKiller.java
 
b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKiller.java
index daaef7d6c03..dd23add68d7 100644
--- 
a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKiller.java
+++ 
b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKiller.java
@@ -45,7 +45,7 @@ import 
org.apache.druid.metadata.SegmentsMetadataManagerConfig;
 import org.apache.druid.metadata.UnusedSegmentKillerConfig;
 import org.apache.druid.query.DruidMetrics;
 import org.apache.druid.segment.loading.DataSegmentKiller;
-import org.apache.druid.timeline.DataSegment;
+import org.apache.druid.server.http.DataSegmentPlus;
 import org.joda.time.DateTime;
 import org.joda.time.Duration;
 import org.joda.time.Interval;
@@ -455,7 +455,7 @@ public class UnusedSegmentsKiller implements OverlordDuty
     }
 
     @Override
-    protected List<DataSegment> fetchNextBatchOfUnusedSegments(TaskToolbox 
toolbox, int nextBatchSize)
+    protected List<DataSegmentPlus> fetchNextBatchOfUnusedSegments(TaskToolbox 
toolbox, int nextBatchSize)
     {
       // Kill only 1000 segments in the batch so that locks are not held for 
very long
       return storageCoordinator.retrieveUnusedSegmentsWithExactInterval(
@@ -466,6 +466,27 @@ public class UnusedSegmentsKiller implements OverlordDuty
       );
     }
 
+    @Override
+    protected Map<String, String> fetchParentIdsForSegments(
+        TaskToolbox toolbox,
+        Map<String, DataSegmentPlus> unusedIdToSegmentPlus
+    )
+    {
+      // No need to make another DB call, the parent IDs have already been 
fetched
+      // in fetchNextBatchOfUnusedSegments
+      final Map<String, String> unusedSegmentIdToParentId = new HashMap<>();
+      for (DataSegmentPlus segment : unusedIdToSegmentPlus.values()) {
+        if (segment.getUpgradedFromSegmentId() != null) {
+          unusedSegmentIdToParentId.put(
+              segment.getDataSegment().getId().toString(),
+              segment.getUpgradedFromSegmentId()
+          );
+        }
+      }
+
+      return unusedSegmentIdToParentId;
+    }
+
     @Override
     protected void logInfo(String message, Object... args)
     {
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/actions/TaskActionTestKit.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/actions/TaskActionTestKit.java
index b6076f26e77..75878ef532f 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/actions/TaskActionTestKit.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/actions/TaskActionTestKit.java
@@ -24,6 +24,7 @@ import com.google.common.base.Preconditions;
 import com.google.common.base.Suppliers;
 import org.apache.druid.indexing.common.TestUtils;
 import org.apache.druid.indexing.common.config.TaskStorageConfig;
+import org.apache.druid.indexing.common.task.Task;
 import org.apache.druid.indexing.overlord.GlobalTaskLockbox;
 import org.apache.druid.indexing.overlord.HeapMemoryTaskStorage;
 import org.apache.druid.indexing.overlord.IndexerMetadataStorageCoordinator;
@@ -53,7 +54,10 @@ import 
org.apache.druid.server.coordinator.simulate.WrappingScheduledExecutorSer
 import org.joda.time.Period;
 import org.junit.rules.ExternalResource;
 
+import java.util.HashMap;
+import java.util.Map;
 import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.Supplier;
 
 public class TaskActionTestKit extends ExternalResource
 {
@@ -73,6 +77,7 @@ public class TaskActionTestKit extends ExternalResource
   private boolean useCentralizedDatasourceSchema = false;
   private boolean batchSegmentAllocation = true;
   private boolean skipSegmentPayloadFetchForAllocation = new 
TaskLockConfig().isBatchAllocationReduceMetadataIO();
+  private Map<Class<? extends TaskAction<?>>, Supplier<?>> taskActionDelegate;
   private AtomicBoolean configFinalized = new AtomicBoolean();
 
   public TaskActionTestKit setUseSegmentMetadataCache(boolean 
useSegmentMetadataCache)
@@ -156,6 +161,36 @@ public class TaskActionTestKit extends ExternalResource
     metadataCachePollExec.finishNextPendingTasks(4);
   }
 
+  /**
+   * Creates a {@link LocalTaskActionClient}. The response for a specific task
+   * action type may be overridden by calling {@link 
#registerDelegateForTaskAction}.
+   */
+  public TaskActionClient createTaskActionClient(Task task)
+  {
+    return new LocalTaskActionClient(task, getTaskActionToolbox())
+    {
+      @Override
+      @SuppressWarnings("unchecked")
+      public <V> V submit(TaskAction<V> taskAction)
+      {
+        final Supplier<?> delegate = 
taskActionDelegate.get(taskAction.getClass());
+        if (delegate == null) {
+          return super.submit(taskAction);
+        } else {
+          return (V) delegate.get();
+        }
+      }
+    };
+  }
+
+  /**
+   * Registers an override action to be performed for task actions of the 
given type.
+   */
+  public <V> void registerDelegateForTaskAction(Class<? extends TaskAction<V>> 
actionType, Supplier<V> function)
+  {
+    taskActionDelegate.put(actionType, function);
+  }
+
   @Override
   public void before()
   {
@@ -234,6 +269,7 @@ public class TaskActionTestKit extends ExternalResource
         supervisorManager,
         objectMapper
     );
+    taskActionDelegate = new HashMap<>();
     testDerbyConnector.createDataSourceTable();
     testDerbyConnector.createUpgradeSegmentsTable();
     testDerbyConnector.createPendingSegmentsTable();
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 7673350e802..eb0f1c0b653 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
@@ -376,7 +376,7 @@ public abstract class IngestionTestBase extends 
InitializedNullHandlingTest
     private final SegmentSchemaMapping segmentSchemaMapping
         = new 
SegmentSchemaMapping(CentralizedDatasourceSchemaConfig.SCHEMA_VERSION);
 
-    private TestLocalTaskActionClient(Task task)
+    public TestLocalTaskActionClient(Task task)
     {
       super(task, taskStorage, getTaskActionToolbox());
     }
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTaskTest.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTaskTest.java
index d1eb4229aa9..a84d28d2f15 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTaskTest.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTaskTest.java
@@ -24,11 +24,15 @@ import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.ImmutableSet;
 import org.apache.druid.error.DruidException;
 import org.apache.druid.error.DruidExceptionMatcher;
+import org.apache.druid.error.ExceptionMatcher;
 import org.apache.druid.indexer.TaskState;
 import org.apache.druid.indexer.report.KillTaskReport;
 import org.apache.druid.indexer.report.TaskReport;
 import org.apache.druid.indexing.common.SegmentLock;
 import org.apache.druid.indexing.common.TaskLockType;
+import 
org.apache.druid.indexing.common.actions.RetrieveUpgradedFromSegmentIdsAction;
+import 
org.apache.druid.indexing.common.actions.RetrieveUpgradedToSegmentIdsAction;
+import org.apache.druid.indexing.common.actions.TaskAction;
 import org.apache.druid.indexing.common.actions.TaskActionClient;
 import org.apache.druid.indexing.common.actions.TimeChunkLockTryAcquireAction;
 import org.apache.druid.indexing.overlord.Segments;
@@ -54,10 +58,12 @@ import org.junit.runner.RunWith;
 import org.junit.runners.Parameterized;
 
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.function.Supplier;
 import java.util.stream.Collectors;
 
 @RunWith(Parameterized.class)
@@ -66,6 +72,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
   private static final String DATA_SOURCE = "wiki";
 
   private TestTaskRunner taskRunner;
+  private Map<Class<? extends TaskAction<?>>, Supplier<Object>> 
taskActionDelegate;
 
   private DataSegment segment1;
   private DataSegment segment2;
@@ -87,6 +94,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
   public void setup()
   {
     taskRunner = new TestTaskRunner();
+    taskActionDelegate = new HashMap<>();
 
     final String version = DateTimes.nowUtc().toString();
     segment1 = newSegment(Intervals.of("2019-01-01/2019-02-01"), 
version).withLoadSpec(ImmutableMap.of("k", 1));
@@ -95,6 +103,25 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     segment4 = newSegment(Intervals.of("2019-04-01/2019-05-01"), 
version).withLoadSpec(ImmutableMap.of("k", 4));
   }
 
+  @Override
+  public TestLocalTaskActionClient createActionClient(Task task)
+  {
+    return new TestLocalTaskActionClient(task)
+    {
+      @Override
+      @SuppressWarnings("unchecked")
+      public <V> V submit(TaskAction<V> taskAction)
+      {
+        final Supplier<Object> delegate = 
taskActionDelegate.get(taskAction.getClass());
+        if (delegate == null) {
+          return super.submit(taskAction);
+        } else {
+          return (V) delegate.get();
+        }
+      }
+    };
+  }
+
   @Test
   public void testKill() throws Exception
   {
@@ -139,7 +166,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     ).containsExactlyInAnyOrder(segment1, segment4);
 
     Assert.assertEquals(
-        new KillTaskReport.Stats(1, 2),
+        new KillTaskReport.Stats(1, 1),
         getReportedStats()
     );
     Assert.assertEquals(ImmutableSet.of(segment3), 
getDataSegmentKiller().getKilledSegments());
@@ -178,7 +205,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     Assert.assertEquals(Collections.emptyList(), observedUnusedSegments);
 
     Assert.assertEquals(
-        new KillTaskReport.Stats(2, 2),
+        new KillTaskReport.Stats(2, 1),
         getReportedStats()
     );
     Assert.assertEquals(ImmutableSet.of(segment1, segment2), 
getDataSegmentKiller().getKilledSegments());
@@ -217,7 +244,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     Assert.assertEquals(Collections.singletonList(segment2), 
observedUnusedSegments);
 
     Assert.assertEquals(
-        new KillTaskReport.Stats(0, 2),
+        new KillTaskReport.Stats(0, 1),
         getReportedStats()
     );
     Assert.assertEquals(Collections.emptySet(), 
getDataSegmentKiller().getKilledSegments());
@@ -262,7 +289,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     ).containsExactlyInAnyOrder(segment1);
 
     Assert.assertEquals(
-        new KillTaskReport.Stats(0, 2),
+        new KillTaskReport.Stats(0, 1),
         getReportedStats()
     );
     Assert.assertEquals(Collections.emptySet(), 
getDataSegmentKiller().getKilledSegments());
@@ -307,7 +334,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     ).containsExactlyInAnyOrder(segment3);
 
     Assert.assertEquals(
-        new KillTaskReport.Stats(0, 2),
+        new KillTaskReport.Stats(0, 1),
         getReportedStats()
     );
     Assert.assertEquals(Collections.emptySet(), 
getDataSegmentKiller().getKilledSegments());
@@ -346,12 +373,101 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
     Assert.assertEquals(ImmutableList.of(), observedUnusedSegments);
 
     Assert.assertEquals(
-        new KillTaskReport.Stats(3, 2),
+        new KillTaskReport.Stats(3, 1),
         getReportedStats()
     );
     Assert.assertEquals(ImmutableSet.of(segment1, segment2, segment3), 
getDataSegmentKiller().getKilledSegments());
   }
 
+  @Test
+  public void 
testTaskFails_andNoSegmentIsDeleted_ifErrorWhileFetchingParentIds()
+  {
+    // Insert some segments and mark them as unused
+    insertUsedSegments(Set.of(segment1, segment2, segment3), Map.of());
+    getMetadataStorageCoordinator().markSegmentAsUnused(segment1.getId());
+    getMetadataStorageCoordinator().markSegmentAsUnused(segment2.getId());
+    getMetadataStorageCoordinator().markSegmentAsUnused(segment3.getId());
+
+    // Make the retrieveUpgradedFromSegmentIds task action fail
+    taskActionDelegate.put(
+        RetrieveUpgradedFromSegmentIdsAction.class,
+        () -> {
+          throw new ISE("Failed to fetch parent IDs");
+        }
+    );
+
+    // Verify that task run fails
+    final KillUnusedSegmentsTask task = new KillUnusedSegmentsTaskBuilder()
+        .dataSource(DATA_SOURCE)
+        .interval(Intervals.ETERNITY)
+        .build();
+
+    MatcherAssert.assertThat(
+        Assert.assertThrows(Exception.class, () -> taskRunner.run(task).get()),
+        ExceptionMatcher.of(Exception.class).expectMessageContains(
+            "Could not retrieve parent segment ids using task 
action[retrieveUpgradedFromSegmentIds]."
+            + " Stopping kill task to avoid data loss in case the segment 
files are shared by other segments."
+        )
+    );
+
+    // Verify that all unused segments are still present in both metadata 
store and deep store
+    final List<DataSegment> observedUnusedSegments =
+        getMetadataStorageCoordinator().retrieveUnusedSegmentsForInterval(
+            DATA_SOURCE,
+            Intervals.ETERNITY,
+            null,
+            null,
+            null
+        );
+    Assert.assertEquals(List.of(segment1, segment2, segment3), 
observedUnusedSegments);
+    Assert.assertEquals(Set.of(), getDataSegmentKiller().getKilledSegments());
+  }
+
+  @Test
+  public void 
testTaskFails_andNoSegmentIsDeleted_ifErrorWhileFetchingChildrenIds()
+  {
+    // Insert some segments and mark them as unused
+    insertUsedSegments(Set.of(segment1, segment2, segment3), Map.of());
+    getMetadataStorageCoordinator().markSegmentAsUnused(segment1.getId());
+    getMetadataStorageCoordinator().markSegmentAsUnused(segment2.getId());
+    getMetadataStorageCoordinator().markSegmentAsUnused(segment3.getId());
+
+    // Make the retrieveUpgradedFromSegmentIds task action fail
+    taskActionDelegate.put(
+        RetrieveUpgradedToSegmentIdsAction.class,
+        () -> {
+          throw new ISE("Failed to fetch children IDs");
+        }
+    );
+
+    // Verify that task run fails
+    final KillUnusedSegmentsTask task = new KillUnusedSegmentsTaskBuilder()
+        .dataSource(DATA_SOURCE)
+        .interval(Intervals.ETERNITY)
+        .build();
+
+    MatcherAssert.assertThat(
+        Assert.assertThrows(Exception.class, () -> taskRunner.run(task).get()),
+        ExceptionMatcher.of(Exception.class).expectMessageContains(
+            "Could not perform task action[retrieveUpgradedToSegmentIds] to 
retrieve"
+            + " segment IDs which share load specs with segments being killed."
+            + " Stopping kill task to avoid data loss in case the segment 
files are shared by other segments."
+        )
+    );
+
+    // Verify that all unused segments are still present in both metadata 
store and deep store
+    final List<DataSegment> observedUnusedSegments =
+        getMetadataStorageCoordinator().retrieveUnusedSegmentsForInterval(
+            DATA_SOURCE,
+            Intervals.ETERNITY,
+            null,
+            null,
+            null
+        );
+    Assert.assertEquals(List.of(segment1, segment2, segment3), 
observedUnusedSegments);
+    Assert.assertEquals(Set.of(), getDataSegmentKiller().getKilledSegments());
+  }
+
   @Test
   public void testKillSegmentsWithVersions() throws Exception
   {
@@ -386,7 +502,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(TaskState.SUCCESS, 
taskRunner.run(task).get().getStatusCode());
     Assert.assertEquals(
-        new KillTaskReport.Stats(4, 3),
+        new KillTaskReport.Stats(4, 2),
         getReportedStats()
     );
 
@@ -435,7 +551,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(TaskState.SUCCESS, 
taskRunner.run(task).get().getStatusCode());
     Assert.assertEquals(
-        new KillTaskReport.Stats(0, 1),
+        new KillTaskReport.Stats(0, 0),
         getReportedStats()
     );
 
@@ -535,7 +651,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(TaskState.SUCCESS, 
taskRunner.run(task).get().getStatusCode());
     Assert.assertEquals(
-        new KillTaskReport.Stats(0, 1),
+        new KillTaskReport.Stats(0, 0),
         getReportedStats()
     );
 
@@ -741,7 +857,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(ImmutableList.of(), observedUnusedSegments);
     Assert.assertEquals(
-        new KillTaskReport.Stats(3, 4),
+        new KillTaskReport.Stats(3, 3),
         getReportedStats()
     );
   }
@@ -827,7 +943,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(ImmutableList.of(segment3), observedUnusedSegments);
     Assert.assertEquals(
-        new KillTaskReport.Stats(2, 3),
+        new KillTaskReport.Stats(2, 2),
         getReportedStats()
     );
 
@@ -851,7 +967,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(ImmutableList.of(), observedUnusedSegments2);
     Assert.assertEquals(
-        new KillTaskReport.Stats(1, 2),
+        new KillTaskReport.Stats(1, 1),
         getReportedStats()
     );
   }
@@ -927,7 +1043,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(ImmutableList.of(segment2, segment3), 
observedUnusedSegments1);
     Assert.assertEquals(
-        new KillTaskReport.Stats(2, 3),
+        new KillTaskReport.Stats(2, 2),
         getReportedStats()
     );
 
@@ -951,7 +1067,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(ImmutableList.of(), observedUnusedSegments2);
     Assert.assertEquals(
-        new KillTaskReport.Stats(2, 3),
+        new KillTaskReport.Stats(2, 2),
         getReportedStats()
     );
   }
@@ -1010,7 +1126,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(TaskState.SUCCESS, 
taskRunner.run(task1).get().getStatusCode());
     Assert.assertEquals(
-        new KillTaskReport.Stats(2, 3),
+        new KillTaskReport.Stats(2, 2),
         getReportedStats()
     );
 
@@ -1035,7 +1151,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(TaskState.SUCCESS, 
taskRunner.run(task2).get().getStatusCode());
     Assert.assertEquals(
-        new KillTaskReport.Stats(1, 2),
+        new KillTaskReport.Stats(1, 1),
         getReportedStats()
     );
 
@@ -1084,7 +1200,7 @@ public class KillUnusedSegmentsTaskTest extends 
IngestionTestBase
 
     Assert.assertEquals(Collections.emptyList(), observedUnusedSegments);
     Assert.assertEquals(
-        new KillTaskReport.Stats(4, 3),
+        new KillTaskReport.Stats(4, 2),
         getReportedStats()
     );
   }
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLifecycleTest.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLifecycleTest.java
index 5da3d2ed754..16149e88a63 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLifecycleTest.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLifecycleTest.java
@@ -846,7 +846,7 @@ public class TaskLifecycleTest extends 
InitializedNullHandlingTest
     Assert.assertEquals("merged statusCode", TaskState.SUCCESS, 
status.getStatusCode());
     Assert.assertEquals("num segments published", 3, 
mdc.getPublished().size());
     Assert.assertEquals("num segments nuked", 3, mdc.getNuked().size());
-    Assert.assertEquals("delete segment batch call count", 2, 
mdc.getDeleteSegmentsCount());
+    Assert.assertEquals("delete segment batch call count", 1, 
mdc.getDeleteSegmentsCount());
     Assert.assertTrue(
         "expected unused segments get killed",
         expectedUnusedSegments.containsAll(mdc.getNuked()) && 
mdc.getNuked().containsAll(
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKillerTest.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKillerTest.java
index 0cc5f6f5f23..7f0d718b0b9 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKillerTest.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/duty/UnusedSegmentsKillerTest.java
@@ -20,7 +20,7 @@
 package org.apache.druid.indexing.overlord.duty;
 
 import org.apache.druid.indexing.common.TaskLockType;
-import org.apache.druid.indexing.common.actions.LocalTaskActionClient;
+import 
org.apache.druid.indexing.common.actions.RetrieveUpgradedToSegmentIdsAction;
 import org.apache.druid.indexing.common.actions.TaskActionTestKit;
 import org.apache.druid.indexing.common.task.NoopTask;
 import org.apache.druid.indexing.common.task.Task;
@@ -29,6 +29,7 @@ import org.apache.druid.indexing.overlord.GlobalTaskLockbox;
 import org.apache.druid.indexing.overlord.IndexerMetadataStorageCoordinator;
 import org.apache.druid.indexing.overlord.TimeChunkLockRequest;
 import org.apache.druid.indexing.test.TestDataSegmentKiller;
+import org.apache.druid.java.util.common.ISE;
 import org.apache.druid.java.util.common.Intervals;
 import org.apache.druid.java.util.common.granularity.Granularities;
 import org.apache.druid.java.util.common.guava.Comparators;
@@ -53,6 +54,7 @@ import org.junit.Rule;
 import org.junit.Test;
 
 import java.util.List;
+import java.util.Map;
 import java.util.Set;
 import java.util.stream.Collectors;
 
@@ -94,7 +96,7 @@ public class UnusedSegmentsKillerTest
             SegmentMetadataCache.UsageMode.ALWAYS,
             killerConfig
         ),
-        task -> new LocalTaskActionClient(task, 
taskActionTestKit.getTaskActionToolbox()),
+        taskActionTestKit::createTaskActionClient,
         storageCoordinator,
         leaderSelector,
         (corePoolSize, nameFormat) -> new 
WrappingScheduledExecutorService(nameFormat, killExecutor, true),
@@ -367,6 +369,32 @@ public class UnusedSegmentsKillerTest
     emitter.verifySum(TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE, 8L);
   }
 
+  @Test
+  public void test_run_isNoop_ifRetrieveUpgradedToSegmentIdsFails()
+  {
+    storageCoordinator.commitSegments(Set.copyOf(WIKI_SEGMENTS_1X10D), null);
+    storageCoordinator.markAllSegmentsAsUnused(TestDataSource.WIKI);
+
+    // Make the retrieveUpgradedFromSegmentIds task action fail
+    taskActionTestKit.registerDelegateForTaskAction(
+        RetrieveUpgradedToSegmentIdsAction.class,
+        () -> {
+          throw new ISE("Failed to fetch children IDs");
+        }
+    );
+
+    leaderSelector.becomeLeader();
+    killer.run();
+
+    // Verify that no unused segment is deleted from metadata store or deep 
store
+    finishQueuedKillJobs();
+    emitter.verifyNotEmitted(TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE);
+    emitter.verifyNotEmitted(TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE);
+
+    // Verify that the task is marked as failed
+    emitter.verifyEmitted("task/run/time", Map.of("taskStatus", "FAILED"), 10);
+  }
+
   @Test
   public void test_run_doesNotKillSegment_ifUpdatedWithinBufferPeriod()
   {
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/test/TestIndexerMetadataStorageCoordinator.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/test/TestIndexerMetadataStorageCoordinator.java
index dbc5a5def5c..41a30024163 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/test/TestIndexerMetadataStorageCoordinator.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/test/TestIndexerMetadataStorageCoordinator.java
@@ -76,7 +76,7 @@ public class TestIndexerMetadataStorageCoordinator implements 
IndexerMetadataSto
   }
 
   @Override
-  public List<DataSegment> retrieveUnusedSegmentsWithExactInterval(
+  public List<DataSegmentPlus> retrieveUnusedSegmentsWithExactInterval(
       String dataSource,
       Interval interval,
       DateTime maxUpdatedTime,
diff --git 
a/server/src/main/java/org/apache/druid/indexing/overlord/IndexerMetadataStorageCoordinator.java
 
b/server/src/main/java/org/apache/druid/indexing/overlord/IndexerMetadataStorageCoordinator.java
index db9b9835ef5..26dcae7cd1f 100644
--- 
a/server/src/main/java/org/apache/druid/indexing/overlord/IndexerMetadataStorageCoordinator.java
+++ 
b/server/src/main/java/org/apache/druid/indexing/overlord/IndexerMetadataStorageCoordinator.java
@@ -161,10 +161,11 @@ public interface IndexerMetadataStorageCoordinator
    * @param maxUpdatedTime Returned segments must have a {@code 
used_status_last_updated}
    *                       which is either null or earlier than this value.
    * @param limit          Maximum number of segments to return.
-   *
    * @return Unsorted list of unused segments that match the given parameters.
+   * The entries in the list are required to have the {@link 
DataSegmentPlus#getDataSegment()}
+   * and {@link DataSegmentPlus#getUpgradedFromSegmentId()} fields populated.
    */
-  List<DataSegment> retrieveUnusedSegmentsWithExactInterval(
+  List<DataSegmentPlus> retrieveUnusedSegmentsWithExactInterval(
       String dataSource,
       Interval interval,
       DateTime maxUpdatedTime,
diff --git 
a/server/src/main/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinator.java
 
b/server/src/main/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinator.java
index 1317fd339a6..fe04a9a0c1d 100644
--- 
a/server/src/main/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinator.java
+++ 
b/server/src/main/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinator.java
@@ -259,7 +259,7 @@ public class IndexerSQLMetadataStorageCoordinator 
implements IndexerMetadataStor
   }
 
   @Override
-  public List<DataSegment> retrieveUnusedSegmentsWithExactInterval(
+  public List<DataSegmentPlus> retrieveUnusedSegmentsWithExactInterval(
       String dataSource,
       Interval interval,
       DateTime maxUpdatedTime,
diff --git 
a/server/src/main/java/org/apache/druid/metadata/SqlSegmentsMetadataQuery.java 
b/server/src/main/java/org/apache/druid/metadata/SqlSegmentsMetadataQuery.java
index 559d3009436..dfeb376204d 100644
--- 
a/server/src/main/java/org/apache/druid/metadata/SqlSegmentsMetadataQuery.java
+++ 
b/server/src/main/java/org/apache/druid/metadata/SqlSegmentsMetadataQuery.java
@@ -1124,7 +1124,7 @@ public class SqlSegmentsMetadataQuery
    *                       which is either null or earlier than this value.
    * @param limit          Maximum number of segments to return
    */
-  public List<DataSegment> retrieveUnusedSegmentsWithExactInterval(
+  public List<DataSegmentPlus> retrieveUnusedSegmentsWithExactInterval(
       String dataSource,
       Interval interval,
       DateTime maxUpdatedTime,
@@ -1132,7 +1132,7 @@ public class SqlSegmentsMetadataQuery
   )
   {
     final String sql = StringUtils.format(
-        "SELECT id, payload FROM %1$s"
+        "SELECT id, payload, upgraded_from_segment_id FROM %1$s"
         + " WHERE dataSource = :dataSource AND used = false"
         + " AND %2$send%2$s = :end AND start = :start"
         + " AND (used_status_last_updated IS NULL OR used_status_last_updated 
<= :maxUpdatedTime)"
@@ -1140,7 +1140,7 @@ public class SqlSegmentsMetadataQuery
         dbTables.getSegmentsTable(), connector.getQuoteString(), 
connector.limitClause(limit)
     );
 
-    final List<DataSegment> segments = connector.inReadOnlyTransaction(
+    final List<DataSegmentPlus> segments = connector.inReadOnlyTransaction(
         (handle, status) ->
             handle.createQuery(sql)
                   .setFetchSize(connector.getStreamingFetchSize())
@@ -1148,7 +1148,7 @@ public class SqlSegmentsMetadataQuery
                   .bind("start", interval.getStart().toString())
                   .bind("end", interval.getEnd().toString())
                   .bind("maxUpdatedTime", maxUpdatedTime.toString())
-                  .map((index, r, ctx) -> mapToSegment(r))
+                  .map((index, r, ctx) -> mapToSegmentPlusUpgradedId(r))
                   .list()
     );
 
@@ -1878,13 +1878,27 @@ public class SqlSegmentsMetadataQuery
     }).iterator();
   }
 
+  /**
+   * Maps the given result set to a {@link DataSegmentPlus} with the segment
+   * payload and the {@code upgradedFromSegmentId} populated (if non-null).
+   */
   @Nullable
-  private DataSegment mapToSegment(ResultSet resultSet)
+  private DataSegmentPlus mapToSegmentPlusUpgradedId(ResultSet resultSet)
   {
     String segmentId = "";
     try {
       segmentId = resultSet.getString("id");
-      return JacksonUtils.readValue(jsonMapper, resultSet.getBytes("payload"), 
DataSegment.class);
+      final String upgradedFromSegmentId = 
resultSet.getString("upgraded_from_segment_id");
+      return new DataSegmentPlus(
+          JacksonUtils.readValue(jsonMapper, resultSet.getBytes("payload"), 
DataSegment.class),
+          null,
+          null,
+          null,
+          null,
+          null,
+          upgradedFromSegmentId,
+          null
+      );
     }
     catch (Throwable t) {
       log.error(t, "Could not read segment with ID[%s]", segmentId);
diff --git 
a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorTest.java
 
b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorTest.java
index 70f4d53cc38..33b65e55f80 100644
--- 
a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorTest.java
+++ 
b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorTest.java
@@ -2217,7 +2217,7 @@ public class IndexerSQLMetadataStorageCoordinatorTest 
extends IndexerSqlMetadata
 
     // Verify that query for exact interval returns the segments
     Assert.assertEquals(
-        List.of(defaultSegment3),
+        List.of(toSegmentPlusUpgradedId(defaultSegment3, null)),
         coordinator.retrieveUnusedSegmentsWithExactInterval(
             dataSource,
             defaultSegment3.getInterval(),
@@ -2228,7 +2228,7 @@ public class IndexerSQLMetadataStorageCoordinatorTest 
extends IndexerSqlMetadata
 
     Assert.assertEquals(defaultSegment.getInterval(), 
defaultSegment2.getInterval());
     Assert.assertEquals(
-        Set.of(defaultSegment, defaultSegment2),
+        Set.of(toSegmentPlusUpgradedId(defaultSegment, null), 
toSegmentPlusUpgradedId(defaultSegment2, null)),
         Set.copyOf(
             coordinator.retrieveUnusedSegmentsWithExactInterval(
                 dataSource,
@@ -4845,4 +4845,9 @@ public class IndexerSQLMetadataStorageCoordinatorTest 
extends IndexerSqlMetadata
         coordinator.retrieveUsedSegmentsForIntervals(dataSource, 
List.of(interval), Segments.ONLY_VISIBLE)
     );
   }
+
+  private DataSegmentPlus toSegmentPlusUpgradedId(DataSegment segment, String 
upgradedFromSegmentId)
+  {
+    return new DataSegmentPlus(segment, null, null, null, null, null, 
upgradedFromSegmentId, null);
+  }
 }


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

Reply via email to