From 626a66d70b2b0030baa6c45849ded21299996a73 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Thu, 23 Jul 2026 22:57:08 +0530 Subject: [PATCH 01/10] Fix bug and improve embedded kill perf --- .../common/task/KillUnusedSegmentsTask.java | 65 ++++++++++++++----- .../overlord/duty/UnusedSegmentsKiller.java | 22 ++++++- ...TestIndexerMetadataStorageCoordinator.java | 2 +- .../IndexerMetadataStorageCoordinator.java | 5 +- .../IndexerSQLMetadataStorageCoordinator.java | 2 +- .../metadata/SqlSegmentsMetadataQuery.java | 26 ++++++-- ...exerSQLMetadataStorageCoordinatorTest.java | 9 ++- 7 files changed, 101 insertions(+), 30 deletions(-) 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 fe58c264ca8b..8ef87598440f 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,6 +48,7 @@ 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; @@ -211,7 +212,7 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception // List unused segments int nextBatchSize = computeNextBatchSize(numSegmentsKilled); @Nullable Integer numTotalBatches = getNumTotalBatches(); - List unusedSegments; + List 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 +237,20 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception break; } - unusedSegments = fetchNextBatchOfUnusedSegments(toolbox, nextBatchSize); + unusedSegmentsPlus = fetchNextBatchOfUnusedSegments(toolbox, nextBatchSize); + if (unusedSegmentsPlus.isEmpty()) { + // No segments eligible for kill + break; + } // Fetch locks each time as a revokal could have occurred in between batches final NavigableMap> taskLockMap = getNonRevokedTaskLockMap(toolbox.getTaskActionClient()); + final Set unusedSegments = unusedSegmentsPlus.stream() + .map(DataSegmentPlus::getDataSegment) + .collect(Collectors.toSet()); + if (!TaskLocks.isLockCoversSegments(taskLockMap, unusedSegments)) { throw new ISE( "Locks[%s] for task[%s] can't cover segments[%s]", @@ -259,28 +268,26 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception // If the segment nuke throws an exception, then the segment cleanup is abandoned. // Determine upgraded segment ids before nuking - final Set segmentIds = unusedSegments.stream() - .map(DataSegment::getId) - .map(SegmentId::toString) - .collect(Collectors.toSet()); final Map upgradedFromSegmentIds = new HashMap<>(); try { upgradedFromSegmentIds.putAll( - taskActionClient.submit( - new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) - ).getUpgradedFromSegmentIds() + fetchParentIdsForSegments(toolbox, unusedSegmentsPlus) ); } catch (Exception e) { - LOG.warn( + // Do not proceed with killing these segments as we cannot be sure if their + // load spec is used by any other segment or not + LOG.error( e, "Could not retrieve parent segment ids using task action[retrieveUpgradedFromSegmentIds]." - + " Overlord may be on an older version." + + " Stopping kill task to avoid deletion of segment files that may be" + + " needed for other segments." ); + break; } // Nuke Segments - taskActionClient.submit(new SegmentNukeAction(new HashSet<>(unusedSegments))); + taskActionClient.submit(new SegmentNukeAction(unusedSegments)); emitMetric(toolbox.getEmitter(), TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size()); // Determine segments to be killed @@ -306,7 +313,7 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception 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 +349,7 @@ int computeNextBatchSize(int numSegmentsKilled) /** * Fetches the next batch of unused segments that are eligible for kill. */ - protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, int nextBatchSize) throws IOException + protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, int nextBatchSize) throws IOException { return toolbox.getTaskActionClient().submit( new RetrieveUnusedSegmentsAction( @@ -352,7 +359,33 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, 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 unusedSegments 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 fetchParentIdsForSegments( + TaskToolbox toolbox, + List unusedSegments + ) throws IOException + { + final Set segmentIds = unusedSegments.stream() + .map(DataSegmentPlus::getDataSegment) + .map(DataSegment::getId) + .map(SegmentId::toString) + .collect(Collectors.toSet()); + return toolbox.getTaskActionClient().submit( + new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) + ).getUpgradedFromSegmentIds(); } /** @@ -387,7 +420,7 @@ private NavigableMap> getNonRevokedTaskLockMap(TaskActi * @return list of segments to kill from deep storage */ private List getKillableSegments( - List unusedSegments, + Set unusedSegments, Map upgradedFromSegmentIds, Set> usedSegmentLoadSpecs, TaskActionClient taskActionClient 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 daaef7d6c03a..7f8efa6882c3 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.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 @@ protected Integer getNumTotalBatches() } @Override - protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, int nextBatchSize) + protected List 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,24 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, ); } + @Override + protected Map fetchParentIdsForSegments(TaskToolbox toolbox, List unusedSegments) + { + // No need to make another DB call, the parent IDs have already been fetched + // in fetchNextBatchOfUnusedSegments + final Map unusedSegmentIdToParentId = new HashMap<>(); + for (DataSegmentPlus segment : unusedSegments) { + 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/test/TestIndexerMetadataStorageCoordinator.java b/indexing-service/src/test/java/org/apache/druid/indexing/test/TestIndexerMetadataStorageCoordinator.java index dbc5a5def5c7..41a300241631 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 List retrieveSomeUnusedSegmentIntervals(String dataSource, int } @Override - public List retrieveUnusedSegmentsWithExactInterval( + public List 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 db9b9835ef5f..77cbd568ba91 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 @@ List retrieveUnusedSegmentsForInterval( * @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 lost are required to have the {@link DataSegmentPlus#getDataSegment()} + * and {@link DataSegmentPlus#getUpgradedFromSegmentId()} fields populated. */ - List retrieveUnusedSegmentsWithExactInterval( + List 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 1317fd339a6a..fe04a9a0c1da 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 List retrieveUnusedSegmentsForInterval( } @Override - public List retrieveUnusedSegmentsWithExactInterval( + public List 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 559d3009436e..dfeb376204d1 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 List retrieveSomeUnusedSegmentIntervals(String dataSource, int * which is either null or earlier than this value. * @param limit Maximum number of segments to return */ - public List retrieveUnusedSegmentsWithExactInterval( + public List retrieveUnusedSegmentsWithExactInterval( String dataSource, Interval interval, DateTime maxUpdatedTime, @@ -1132,7 +1132,7 @@ public List retrieveUnusedSegmentsWithExactInterval( ) { 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 List retrieveUnusedSegmentsWithExactInterval( dbTables.getSegmentsTable(), connector.getQuoteString(), connector.limitClause(limit) ); - final List segments = connector.inReadOnlyTransaction( + final List segments = connector.inReadOnlyTransaction( (handle, status) -> handle.createQuery(sql) .setFetchSize(connector.getStreamingFetchSize()) @@ -1148,7 +1148,7 @@ public List retrieveUnusedSegmentsWithExactInterval( .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 @@ private ResultIterator getDataSegmentPlusResultIterator( }).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 70f4d53cc38a..33b65e55f800 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 void testRetrieveUnusedSegmentsWithExactInterval() // 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 void testRetrieveUnusedSegmentsWithExactInterval() 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 @@ private void verifyIntervalHasVisibleSegments( 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); + } } From 392b190db29b4305c4400ce09820e880edf8af35 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Thu, 23 Jul 2026 23:12:27 +0530 Subject: [PATCH 02/10] minor fixes --- .../common/task/KillUnusedSegmentsTask.java | 21 ++++++++++--------- 1 file changed, 11 insertions(+), 10 deletions(-) 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 8ef87598440f..b0a14f980145 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 @@ -52,7 +52,6 @@ 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; @@ -239,7 +238,7 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception unusedSegmentsPlus = fetchNextBatchOfUnusedSegments(toolbox, nextBatchSize); if (unusedSegmentsPlus.isEmpty()) { - // No segments eligible for kill + // No more segments eligible for kill, do not proceed further break; } @@ -276,12 +275,15 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception } catch (Exception e) { // Do not proceed with killing these segments as we cannot be sure if their - // load spec is used by any other segment or not + // 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. LOG.error( e, "Could not retrieve parent segment ids using task action[retrieveUpgradedFromSegmentIds]." - + " Stopping kill task to avoid deletion of segment files that may be" - + " needed for other segments." + + " Stopping kill task to avoid data loss in case the segment files" + + " are shared by other segments." ); break; } @@ -378,11 +380,10 @@ protected Map fetchParentIdsForSegments( List unusedSegments ) throws IOException { - final Set segmentIds = unusedSegments.stream() - .map(DataSegmentPlus::getDataSegment) - .map(DataSegment::getId) - .map(SegmentId::toString) - .collect(Collectors.toSet()); + final Set segmentIds = unusedSegments.stream().map( + s -> s.getDataSegment().getId().toString() + ).collect(Collectors.toSet()); + return toolbox.getTaskActionClient().submit( new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) ).getUpgradedFromSegmentIds(); From 8f6c8200acece0301701eb7acf842231a42a5311 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Fri, 24 Jul 2026 00:36:49 +0530 Subject: [PATCH 03/10] More fixes --- .../common/task/KillUnusedSegmentsTask.java | 129 ++++++++++-------- .../overlord/duty/UnusedSegmentsKiller.java | 5 +- 2 files changed, 76 insertions(+), 58 deletions(-) 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 b0a14f980145..73ca586026ee 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 @@ -24,6 +24,7 @@ import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.annotation.JsonProperty; import com.google.common.annotations.VisibleForTesting; +import com.google.common.base.Optional; import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableSet; import org.apache.druid.client.indexing.ClientKillUnusedSegmentsTaskQuery; @@ -259,58 +260,50 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception ); } - // 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 Map upgradedFromSegmentIds = new HashMap<>(); - try { - upgradedFromSegmentIds.putAll( - fetchParentIdsForSegments(toolbox, unusedSegmentsPlus) - ); - } - 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. - LOG.error( - 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." - ); + // Kill segments - order of steps 1, 2, 3, 4 must remain the same + + // 1. Determine parent segment ids of killable unused segments + final Optional> upgradedFromSegmentIds + = fetchParentIdsForSegments(toolbox, unusedSegmentsPlus); + if (!upgradedFromSegmentIds.isPresent()) { + // Do not proceed further as we do not know for sure if load specs are shared or not break; } - // Nuke Segments - taskActionClient.submit(new SegmentNukeAction(unusedSegments)); - emitMetric(toolbox.getEmitter(), TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size()); - - // Determine segments to be killed - final List segmentsToBeKilled - = getKillableSegments(unusedSegments, upgradedFromSegmentIds, usedSegmentLoadSpecs, taskActionClient); + // 2. Identify killable segments whose load specs are not shared with any other segment + final List segmentsToKillFromDeepStore = getKillableSegments( + unusedSegments, + upgradedFromSegmentIds.get(), + usedSegmentLoadSpecs, + taskActionClient + ); + if (segmentsToKillFromDeepStore.isEmpty()) { + // Do not proceed further as we do not know for sure if load specs are shared or not + break; + } + // 2a. Track segments that cannot be removed from deep store yet final Set 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 used 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, but only if we know for sure if their + // load specs are shared by other segments or not + taskActionClient.submit(new SegmentNukeAction(unusedSegments)); + emitMetric(toolbox.getEmitter(), TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.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()); @@ -371,22 +364,42 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolb * Fetches the parent IDs (if any) for the given unused segments. * * @param unusedSegments 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. + * @return Optional containing 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. Empty optional if an error occurred. */ - protected Map fetchParentIdsForSegments( + protected Optional> fetchParentIdsForSegments( TaskToolbox toolbox, List unusedSegments - ) throws IOException + ) { - final Set segmentIds = unusedSegments.stream().map( - s -> s.getDataSegment().getId().toString() - ).collect(Collectors.toSet()); + try { + final Set segmentIds = unusedSegments.stream().map( + s -> s.getDataSegment().getId().toString() + ).collect(Collectors.toSet()); - return toolbox.getTaskActionClient().submit( - new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) - ).getUpgradedFromSegmentIds(); + final Map segmentIdToParentId = toolbox.getTaskActionClient().submit( + new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) + ).getUpgradedFromSegmentIds(); + + return segmentIdToParentId == null + ? Optional.of(Map.of()) + : Optional.of(segmentIdToParentId); + } + 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. + LOG.error( + 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." + ); + return Optional.absent(); + } } /** @@ -427,7 +440,6 @@ private List getKillableSegments( TaskActionClient taskActionClient ) { - // Determine parentId for each unused segment final Map> parentIdToUnusedSegments = new HashMap<>(); for (DataSegment segment : unusedSegments) { @@ -457,11 +469,16 @@ private List getKillableSegments( } } 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 + LOG.error( 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." ); + return List.of(); } // Filter using the used segment load specs as segment upgrades predate the above task action @@ -483,7 +500,7 @@ private boolean isSegmentLoadSpecPresentIn( { 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 7f8efa6882c3..a378172b3086 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 @@ -19,6 +19,7 @@ package org.apache.druid.indexing.overlord.duty; +import com.google.common.base.Optional; import com.google.common.collect.Ordering; import com.google.inject.Inject; import org.apache.druid.client.indexing.IndexingService; @@ -467,7 +468,7 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolb } @Override - protected Map fetchParentIdsForSegments(TaskToolbox toolbox, List unusedSegments) + protected Optional> fetchParentIdsForSegments(TaskToolbox toolbox, List unusedSegments) { // No need to make another DB call, the parent IDs have already been fetched // in fetchNextBatchOfUnusedSegments @@ -481,7 +482,7 @@ protected Map fetchParentIdsForSegments(TaskToolbox toolbox, Lis } } - return unusedSegmentIdToParentId; + return Optional.of(unusedSegmentIdToParentId); } @Override From 4249ea96bcee56afd6eb4a968ed9e7408b595011 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Fri, 24 Jul 2026 01:21:33 +0530 Subject: [PATCH 04/10] Exclude self --- .../common/task/KillUnusedSegmentsTask.java | 55 ++++++++----------- .../overlord/duty/UnusedSegmentsKiller.java | 5 +- 2 files changed, 26 insertions(+), 34 deletions(-) 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 73ca586026ee..c91401628889 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 @@ -24,7 +24,6 @@ import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.annotation.JsonProperty; import com.google.common.annotations.VisibleForTesting; -import com.google.common.base.Optional; import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableSet; import org.apache.druid.client.indexing.ClientKillUnusedSegmentsTaskQuery; @@ -53,7 +52,6 @@ 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.utils.CollectionUtils; import org.joda.time.DateTime; import org.joda.time.Interval; @@ -263,24 +261,16 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception // Kill segments - order of steps 1, 2, 3, 4 must remain the same // 1. Determine parent segment ids of killable unused segments - final Optional> upgradedFromSegmentIds + final Map upgradedFromSegmentIds = fetchParentIdsForSegments(toolbox, unusedSegmentsPlus); - if (!upgradedFromSegmentIds.isPresent()) { - // Do not proceed further as we do not know for sure if load specs are shared or not - break; - } // 2. Identify killable segments whose load specs are not shared with any other segment final List segmentsToKillFromDeepStore = getKillableSegments( unusedSegments, - upgradedFromSegmentIds.get(), + upgradedFromSegmentIds, usedSegmentLoadSpecs, taskActionClient ); - if (segmentsToKillFromDeepStore.isEmpty()) { - // Do not proceed further as we do not know for sure if load specs are shared or not - break; - } // 2a. Track segments that cannot be removed from deep store yet final Set segmentsNotKilled = new HashSet<>(unusedSegments); @@ -288,13 +278,16 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception if (!segmentsNotKilled.isEmpty()) { LOG.warn( "Skipping kill of [%d] segments of datasource[%s] from deep storage" - + " as their load specs are used by other segments.", + + " as their load specs are shared by other segments.", segmentsNotKilled.size(), getDataSource() ); } + if (segmentsToKillFromDeepStore.isEmpty()) { + // Do not proceed further as we will always keep getting the same batch of segments + break; + } - // 3. Nuke all eligible unused segments, but only if we know for sure if their - // load specs are shared by other segments or not + // 3. Nuke all eligible unused segments taskActionClient.submit(new SegmentNukeAction(unusedSegments)); emitMetric(toolbox.getEmitter(), TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size()); @@ -364,11 +357,11 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolb * Fetches the parent IDs (if any) for the given unused segments. * * @param unusedSegments Unused segments whose parent IDs need to be fetched - * @return Optional containing Map from segment ID to the segment ID from which + * @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. Empty optional if an error occurred. + * it does not have an entry in the map. */ - protected Optional> fetchParentIdsForSegments( + protected Map fetchParentIdsForSegments( TaskToolbox toolbox, List unusedSegments ) @@ -378,13 +371,9 @@ protected Optional> fetchParentIdsForSegments( s -> s.getDataSegment().getId().toString() ).collect(Collectors.toSet()); - final Map segmentIdToParentId = toolbox.getTaskActionClient().submit( + return toolbox.getTaskActionClient().submit( new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) ).getUpgradedFromSegmentIds(); - - return segmentIdToParentId == null - ? Optional.of(Map.of()) - : Optional.of(segmentIdToParentId); } catch (Exception e) { // Do not proceed with killing these segments as we cannot be sure if their @@ -392,13 +381,12 @@ protected Optional> fetchParentIdsForSegments( // 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. - LOG.error( + 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." ); - return Optional.absent(); } } @@ -440,7 +428,12 @@ private List getKillableSegments( TaskActionClient taskActionClient ) { - // Determine parentId for each unused segment + // Unused segment IDs being killed + final Set segmentIdsBeingKilled = unusedSegments.stream() + .map(s -> s.getId().toString()) + .collect(Collectors.toSet()); + + // Determine parentId (or self, if no parent) for each unused segment final Map> parentIdToUnusedSegments = new HashMap<>(); for (DataSegment segment : unusedSegments) { final String segmentId = segment.getId().toString(); @@ -457,10 +450,11 @@ private List getKillableSegments( ); 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 (!segmentIdsBeingKilled.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); @@ -471,14 +465,13 @@ private List getKillableSegments( catch (Exception e) { // Do not proceed with the kill of any segment as we cannot be sure if their // load specs are shared by any other segment - LOG.error( + throw new ISE( e, "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." ); - return List.of(); } // Filter using the used segment load specs as segment upgrades predate the above task action 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 a378172b3086..7f8efa6882c3 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 @@ -19,7 +19,6 @@ package org.apache.druid.indexing.overlord.duty; -import com.google.common.base.Optional; import com.google.common.collect.Ordering; import com.google.inject.Inject; import org.apache.druid.client.indexing.IndexingService; @@ -468,7 +467,7 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolb } @Override - protected Optional> fetchParentIdsForSegments(TaskToolbox toolbox, List unusedSegments) + protected Map fetchParentIdsForSegments(TaskToolbox toolbox, List unusedSegments) { // No need to make another DB call, the parent IDs have already been fetched // in fetchNextBatchOfUnusedSegments @@ -482,7 +481,7 @@ protected Optional> fetchParentIdsForSegments(TaskToolbox to } } - return Optional.of(unusedSegmentIdToParentId); + return unusedSegmentIdToParentId; } @Override From 02612d9d2be7457d37fcaeaa707fe5ec9011b515 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Fri, 24 Jul 2026 14:00:42 +0530 Subject: [PATCH 05/10] Fix tests --- .../common/task/KillUnusedSegmentsTask.java | 4 --- .../task/KillUnusedSegmentsTaskTest.java | 34 +++++++++---------- .../indexing/overlord/TaskLifecycleTest.java | 2 +- 3 files changed, 18 insertions(+), 22 deletions(-) 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 c91401628889..6e7b50b83f78 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 @@ -282,10 +282,6 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception segmentsNotKilled.size(), getDataSource() ); } - if (segmentsToKillFromDeepStore.isEmpty()) { - // Do not proceed further as we will always keep getting the same batch of segments - break; - } // 3. Nuke all eligible unused segments taskActionClient.submit(new SegmentNukeAction(unusedSegments)); 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 d1eb4229aa91..8899f14199af 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 @@ -139,7 +139,7 @@ public void testKill() throws Exception ).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 +178,7 @@ public void testKillSegmentsDeleteUnreferencedSiblings() throws Exception 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 +217,7 @@ public void testKillSegmentsDoNotDeleteReferencedSibling() throws Exception 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 +262,7 @@ public void testKillSegmentsDoNotDeleteParentWithReferencedChildren() throws Exc ).containsExactlyInAnyOrder(segment1); Assert.assertEquals( - new KillTaskReport.Stats(0, 2), + new KillTaskReport.Stats(0, 1), getReportedStats() ); Assert.assertEquals(Collections.emptySet(), getDataSegmentKiller().getKilledSegments()); @@ -307,7 +307,7 @@ public void testKillSegmentsDoNotDeleteChildrenWithReferencedParent() throws Exc ).containsExactlyInAnyOrder(segment3); Assert.assertEquals( - new KillTaskReport.Stats(0, 2), + new KillTaskReport.Stats(0, 1), getReportedStats() ); Assert.assertEquals(Collections.emptySet(), getDataSegmentKiller().getKilledSegments()); @@ -346,7 +346,7 @@ public void testKillSegmentsDeleteChildrenAndParent() throws Exception 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()); @@ -386,7 +386,7 @@ public void testKillSegmentsWithVersions() throws Exception Assert.assertEquals(TaskState.SUCCESS, taskRunner.run(task).get().getStatusCode()); Assert.assertEquals( - new KillTaskReport.Stats(4, 3), + new KillTaskReport.Stats(4, 2), getReportedStats() ); @@ -435,7 +435,7 @@ public void testKillSegmentsWithEmptyVersions() throws Exception Assert.assertEquals(TaskState.SUCCESS, taskRunner.run(task).get().getStatusCode()); Assert.assertEquals( - new KillTaskReport.Stats(0, 1), + new KillTaskReport.Stats(0, 0), getReportedStats() ); @@ -535,7 +535,7 @@ public void testKillWithNonExistentVersion() throws Exception Assert.assertEquals(TaskState.SUCCESS, taskRunner.run(task).get().getStatusCode()); Assert.assertEquals( - new KillTaskReport.Stats(0, 1), + new KillTaskReport.Stats(0, 0), getReportedStats() ); @@ -741,7 +741,7 @@ public void testKillMultipleUnusedSegmentsWithNullMaxUsedStatusLastUpdatedTime() Assert.assertEquals(ImmutableList.of(), observedUnusedSegments); Assert.assertEquals( - new KillTaskReport.Stats(3, 4), + new KillTaskReport.Stats(3, 3), getReportedStats() ); } @@ -827,7 +827,7 @@ public void testKillMultipleUnusedSegmentsWithDifferentMaxUsedStatusLastUpdatedT Assert.assertEquals(ImmutableList.of(segment3), observedUnusedSegments); Assert.assertEquals( - new KillTaskReport.Stats(2, 3), + new KillTaskReport.Stats(2, 2), getReportedStats() ); @@ -851,7 +851,7 @@ public void testKillMultipleUnusedSegmentsWithDifferentMaxUsedStatusLastUpdatedT Assert.assertEquals(ImmutableList.of(), observedUnusedSegments2); Assert.assertEquals( - new KillTaskReport.Stats(1, 2), + new KillTaskReport.Stats(1, 1), getReportedStats() ); } @@ -927,7 +927,7 @@ public void testKillMultipleUnusedSegmentsWithDifferentMaxUsedStatusLastUpdatedT Assert.assertEquals(ImmutableList.of(segment2, segment3), observedUnusedSegments1); Assert.assertEquals( - new KillTaskReport.Stats(2, 3), + new KillTaskReport.Stats(2, 2), getReportedStats() ); @@ -951,7 +951,7 @@ public void testKillMultipleUnusedSegmentsWithDifferentMaxUsedStatusLastUpdatedT Assert.assertEquals(ImmutableList.of(), observedUnusedSegments2); Assert.assertEquals( - new KillTaskReport.Stats(2, 3), + new KillTaskReport.Stats(2, 2), getReportedStats() ); } @@ -1010,7 +1010,7 @@ public void testKillMultipleUnusedSegmentsWithVersionAndDifferentLastUpdatedTime Assert.assertEquals(TaskState.SUCCESS, taskRunner.run(task1).get().getStatusCode()); Assert.assertEquals( - new KillTaskReport.Stats(2, 3), + new KillTaskReport.Stats(2, 2), getReportedStats() ); @@ -1035,7 +1035,7 @@ public void testKillMultipleUnusedSegmentsWithVersionAndDifferentLastUpdatedTime Assert.assertEquals(TaskState.SUCCESS, taskRunner.run(task2).get().getStatusCode()); Assert.assertEquals( - new KillTaskReport.Stats(1, 2), + new KillTaskReport.Stats(1, 1), getReportedStats() ); @@ -1084,7 +1084,7 @@ public void testKillBatchSizeThree() throws Exception 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 5da3d2ed754a..16149e88a630 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 DataSegment apply(String input) 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( From cb459852097f297eeb479cbf3521e589845b1e4b Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Fri, 24 Jul 2026 15:39:40 +0530 Subject: [PATCH 06/10] Add some tests --- .../common/task/IngestionTestBase.java | 2 +- .../task/KillUnusedSegmentsTaskTest.java | 116 ++++++++++++++++++ 2 files changed, 117 insertions(+), 1 deletion(-) 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 7673350e8029..eb0f1c0b653d 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 class TestLocalTaskActionClient extends CountingLocalTaskActionClientForT 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 8899f14199af..a84d28d2f155 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.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.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>, Supplier> taskActionDelegate; private DataSegment segment1; private DataSegment segment2; @@ -87,6 +94,7 @@ public KillUnusedSegmentsTaskTest(boolean useSegmentMetadataCache) 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 void setup() 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 submit(TaskAction taskAction) + { + final Supplier delegate = taskActionDelegate.get(taskAction.getClass()); + if (delegate == null) { + return super.submit(taskAction); + } else { + return (V) delegate.get(); + } + } + }; + } + @Test public void testKill() throws Exception { @@ -352,6 +379,95 @@ public void testKillSegmentsDeleteChildrenAndParent() throws Exception 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 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 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 { From 8cadbbeee2d8bdc1074a9dea5340b31d7c08af7d Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Fri, 24 Jul 2026 16:19:12 +0530 Subject: [PATCH 07/10] Add tests for UnusedSegmentsKiller --- .../common/actions/TaskActionTestKit.java | 36 +++++++++++++++++++ .../duty/UnusedSegmentsKillerTest.java | 32 +++++++++++++++-- 2 files changed, 66 insertions(+), 2 deletions(-) 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 b6076f26e771..75878ef532f3 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.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.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>, Supplier> taskActionDelegate; private AtomicBoolean configFinalized = new AtomicBoolean(); public TaskActionTestKit setUseSegmentMetadataCache(boolean useSegmentMetadataCache) @@ -156,6 +161,36 @@ public void syncSegmentMetadataCache() 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 submit(TaskAction 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 void registerDelegateForTaskAction(Class> actionType, Supplier function) + { + taskActionDelegate.put(actionType, function); + } + @Override public void before() { @@ -234,6 +269,7 @@ public boolean isBatchAllocationReduceMetadataIO() supervisorManager, objectMapper ); + taskActionDelegate = new HashMap<>(); testDerbyConnector.createDataSourceTable(); testDerbyConnector.createUpgradeSegmentsTable(); testDerbyConnector.createPendingSegmentsTable(); 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 0cc5f6f5f232..7f0d718b0b9a 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.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.Test; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.stream.Collectors; @@ -94,7 +96,7 @@ private void initKiller() 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 void test_run_doesNotDeleteSegmentFiles_ifLoadSpecIsUsedByAnotherSegment( 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() { From 0f2ee05207d27f690a3c6aea22bdf10b017e8b84 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Fri, 24 Jul 2026 21:35:27 +0530 Subject: [PATCH 08/10] Add note in javadoc --- .../druid/indexing/common/task/KillUnusedSegmentsTask.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) 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 6e7b50b83f78..1669cb015b70 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 @@ -69,7 +69,6 @@ import java.util.stream.Collectors; /** - *

* 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. @@ -86,6 +85,10 @@ *

  • Filter the set of unreferenced segments using load specs from the set of used segments.
  • *
  • Kill the filtered set of segments from deep storage.
  • * + * 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 { From 0dd7e69b759c2237d876745d83b4147730839eaf Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Sat, 25 Jul 2026 16:16:29 +0530 Subject: [PATCH 09/10] fix docs, javadoc --- docs/data-management/delete.md | 8 ++++++-- .../overlord/IndexerMetadataStorageCoordinator.java | 2 +- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/docs/data-management/delete.md b/docs/data-management/delete.md index cf571566edec..799e8b4b8b91 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/server/src/main/java/org/apache/druid/indexing/overlord/IndexerMetadataStorageCoordinator.java b/server/src/main/java/org/apache/druid/indexing/overlord/IndexerMetadataStorageCoordinator.java index 77cbd568ba91..26dcae7cd1f4 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 @@ -162,7 +162,7 @@ List retrieveUnusedSegmentsForInterval( * 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 lost are required to have the {@link DataSegmentPlus#getDataSegment()} + * The entries in the list are required to have the {@link DataSegmentPlus#getDataSegment()} * and {@link DataSegmentPlus#getUpgradedFromSegmentId()} fields populated. */ List retrieveUnusedSegmentsWithExactInterval( From 3728d11b7636d973040665250348ed8c2470f619 Mon Sep 17 00:00:00 2001 From: Kashif Faraz Date: Sat, 25 Jul 2026 16:30:54 +0530 Subject: [PATCH 10/10] Minor fixes --- .../common/task/KillUnusedSegmentsTask.java | 43 +++++++++---------- .../overlord/duty/UnusedSegmentsKiller.java | 7 ++- 2 files changed, 26 insertions(+), 24 deletions(-) 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 1669cb015b70..21b9e3f6a799 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 @@ -52,6 +52,7 @@ 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.utils.CollectionUtils; import org.joda.time.DateTime; import org.joda.time.Interval; @@ -66,6 +67,7 @@ import java.util.NavigableMap; import java.util.Set; import java.util.TreeMap; +import java.util.function.Function; import java.util.stream.Collectors; /** @@ -249,8 +251,13 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception = getNonRevokedTaskLockMap(toolbox.getTaskActionClient()); final Set unusedSegments = unusedSegmentsPlus.stream() - .map(DataSegmentPlus::getDataSegment) - .collect(Collectors.toSet()); + .map(DataSegmentPlus::getDataSegment) + .collect(Collectors.toSet()); + final Map unusedIdToSegmentPlus = CollectionUtils.toMap( + unusedSegmentsPlus, + segment -> segment.getDataSegment().getId().toString(), + Function.identity() + ); if (!TaskLocks.isLockCoversSegments(taskLockMap, unusedSegments)) { throw new ISE( @@ -265,11 +272,11 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception // 1. Determine parent segment ids of killable unused segments final Map upgradedFromSegmentIds - = fetchParentIdsForSegments(toolbox, unusedSegmentsPlus); + = fetchParentIdsForSegments(toolbox, unusedIdToSegmentPlus); // 2. Identify killable segments whose load specs are not shared with any other segment final List segmentsToKillFromDeepStore = getKillableSegments( - unusedSegments, + unusedIdToSegmentPlus, upgradedFromSegmentIds, usedSegmentLoadSpecs, taskActionClient @@ -288,7 +295,7 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception // 3. Nuke all eligible unused segments taskActionClient.submit(new SegmentNukeAction(unusedSegments)); - emitMetric(toolbox.getEmitter(), TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size()); + 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); @@ -355,23 +362,20 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolb /** * Fetches the parent IDs (if any) for the given unused segments. * - * @param unusedSegments Unused segments whose parent IDs need to be fetched + * @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 fetchParentIdsForSegments( TaskToolbox toolbox, - List unusedSegments + Map unusedIdToSegmentPlus ) { try { - final Set segmentIds = unusedSegments.stream().map( - s -> s.getDataSegment().getId().toString() - ).collect(Collectors.toSet()); - return toolbox.getTaskActionClient().submit( - new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds) + new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), unusedIdToSegmentPlus.keySet()) ).getUpgradedFromSegmentIds(); } catch (Exception e) { @@ -421,25 +425,20 @@ private NavigableMap> getNonRevokedTaskLockMap(TaskActi * @return list of segments to kill from deep storage */ private List getKillableSegments( - Set unusedSegments, + Map unusedSegments, Map upgradedFromSegmentIds, Set> usedSegmentLoadSpecs, TaskActionClient taskActionClient ) { - // Unused segment IDs being killed - final Set segmentIdsBeingKilled = unusedSegments.stream() - .map(s -> s.getId().toString()) - .collect(Collectors.toSet()); - // Determine parentId (or self, if no parent) for each unused segment final Map> parentIdToUnusedSegments = new HashMap<>(); - for (DataSegment segment : unusedSegments) { - final String segmentId = segment.getId().toString(); + for (Map.Entry 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 @@ -449,7 +448,7 @@ private List getKillableSegments( ); if (response != null && response.getUpgradedToSegmentIds() != null) { response.getUpgradedToSegmentIds().forEach((parent, children) -> { - if (!segmentIdsBeingKilled.containsAll(children)) { + 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( 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 7f8efa6882c3..dd23add68d7f 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 @@ -467,12 +467,15 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolb } @Override - protected Map fetchParentIdsForSegments(TaskToolbox toolbox, List unusedSegments) + protected Map fetchParentIdsForSegments( + TaskToolbox toolbox, + Map unusedIdToSegmentPlus + ) { // No need to make another DB call, the parent IDs have already been fetched // in fetchNextBatchOfUnusedSegments final Map unusedSegmentIdToParentId = new HashMap<>(); - for (DataSegmentPlus segment : unusedSegments) { + for (DataSegmentPlus segment : unusedIdToSegmentPlus.values()) { if (segment.getUpgradedFromSegmentId() != null) { unusedSegmentIdToParentId.put( segment.getDataSegment().getId().toString(),