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/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..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 @@ -48,10 +48,10 @@ 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.NavigableMap; import java.util.Set; import java.util.TreeMap; +import java.util.function.Function; 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. @@ -87,6 +87,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 { @@ -211,7 +215,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 +240,25 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception 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> taskLockMap = getNonRevokedTaskLockMap(toolbox.getTaskActionClient()); + final Set unusedSegments = unusedSegmentsPlus.stream() + .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( "Locks[%s] for task[%s] can't cover segments[%s]", @@ -251,62 +268,46 @@ 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 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() - ); - } - 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 upgradedFromSegmentIds + = fetchParentIdsForSegments(toolbox, unusedIdToSegmentPlus); - // 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( + unusedIdToSegmentPlus, + upgradedFromSegmentIds, + usedSegmentLoadSpecs, + taskActionClient + ); + // 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 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 @@ 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 +353,44 @@ 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 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, + Map 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 @@ private NavigableMap> getNonRevokedTaskLockMap(TaskActi * @return list of segments to kill from deep storage */ private List getKillableSegments( - List unusedSegments, + Map unusedSegments, Map upgradedFromSegmentIds, Set> usedSegmentLoadSpecs, TaskActionClient taskActionClient ) { - - // Determine parentId for each unused segment + // 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 @@ -411,10 +448,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 (!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 @@ 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 + 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 @@ 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 daaef7d6c03a..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 @@ -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,27 @@ protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, ); } + @Override + 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 : 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 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/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 d1eb4229aa91..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 { @@ -139,7 +166,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 +205,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 +244,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 +289,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 +334,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,12 +373,101 @@ 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()); } + @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 { @@ -386,7 +502,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 +551,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 +651,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 +857,7 @@ public void testKillMultipleUnusedSegmentsWithNullMaxUsedStatusLastUpdatedTime() Assert.assertEquals(ImmutableList.of(), observedUnusedSegments); Assert.assertEquals( - new KillTaskReport.Stats(3, 4), + new KillTaskReport.Stats(3, 3), getReportedStats() ); } @@ -827,7 +943,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 +967,7 @@ public void testKillMultipleUnusedSegmentsWithDifferentMaxUsedStatusLastUpdatedT Assert.assertEquals(ImmutableList.of(), observedUnusedSegments2); Assert.assertEquals( - new KillTaskReport.Stats(1, 2), + new KillTaskReport.Stats(1, 1), getReportedStats() ); } @@ -927,7 +1043,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 +1067,7 @@ public void testKillMultipleUnusedSegmentsWithDifferentMaxUsedStatusLastUpdatedT Assert.assertEquals(ImmutableList.of(), observedUnusedSegments2); Assert.assertEquals( - new KillTaskReport.Stats(2, 3), + new KillTaskReport.Stats(2, 2), getReportedStats() ); } @@ -1010,7 +1126,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 +1151,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 +1200,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( 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() { 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..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 @@ -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 list 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); + } }