abhishekrb19 commented on code in PR #19737:
URL: https://github.com/apache/druid/pull/19737#discussion_r3643738051


##########
indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java:
##########
@@ -251,62 +258,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 Set<String> segmentIds = unusedSegments.stream()
-                                                   .map(DataSegment::getId)
-                                                   .map(SegmentId::toString)
-                                                   
.collect(Collectors.toSet());
-      final Map<String, String> upgradedFromSegmentIds = new HashMap<>();
-      try {
-        upgradedFromSegmentIds.putAll(
-            taskActionClient.submit(
-                new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), 
segmentIds)
-            ).getUpgradedFromSegmentIds()
-        );
-      }
-      catch (Exception e) {
-        LOG.warn(
-            e,
-            "Could not retrieve parent segment ids using task 
action[retrieveUpgradedFromSegmentIds]."
-            + " Overlord may be on an older version."
-        );
-      }
+      // Kill segments - order of steps 1, 2, 3, 4 must remain the same
 
-      // Nuke Segments
-      taskActionClient.submit(new SegmentNukeAction(new 
HashSet<>(unusedSegments)));
-      emitMetric(toolbox.getEmitter(), 
TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size());
+      // 1. Determine parent segment ids of killable unused segments
+      final Map<String, String> upgradedFromSegmentIds
+          = fetchParentIdsForSegments(toolbox, unusedSegmentsPlus);
 
-      // Determine segments to be killed
-      final List<DataSegment> segmentsToBeKilled
-          = getKillableSegments(unusedSegments, upgradedFromSegmentIds, 
usedSegmentLoadSpecs, taskActionClient);
+      // 2. Identify killable segments whose load specs are not shared with 
any other segment
+      final List<DataSegment> segmentsToKillFromDeepStore = 
getKillableSegments(
+          unusedSegments,
+          upgradedFromSegmentIds,
+          usedSegmentLoadSpecs,
+          taskActionClient
+      );
 
+      // 2a. Track segments that cannot be removed from deep store yet
       final Set<DataSegment> segmentsNotKilled = new HashSet<>(unusedSegments);
-      segmentsToBeKilled.forEach(segmentsNotKilled::remove);
-
+      segmentsToKillFromDeepStore.forEach(segmentsNotKilled::remove);
       if (!segmentsNotKilled.isEmpty()) {
         LOG.warn(
-            "Skipping kill of [%d] segments from deep storage as their load 
specs are used by other segments.",
-            segmentsNotKilled.size()
+            "Skipping kill of [%d] segments of datasource[%s] from deep 
storage"

Review Comment:
   Would it help to log the first N segment IDs or just a small list of 
intervals?



##########
indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java:
##########
@@ -251,62 +258,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 Set<String> segmentIds = unusedSegments.stream()
-                                                   .map(DataSegment::getId)
-                                                   .map(SegmentId::toString)
-                                                   
.collect(Collectors.toSet());
-      final Map<String, String> upgradedFromSegmentIds = new HashMap<>();
-      try {
-        upgradedFromSegmentIds.putAll(
-            taskActionClient.submit(
-                new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), 
segmentIds)
-            ).getUpgradedFromSegmentIds()
-        );
-      }
-      catch (Exception e) {
-        LOG.warn(
-            e,
-            "Could not retrieve parent segment ids using task 
action[retrieveUpgradedFromSegmentIds]."
-            + " Overlord may be on an older version."
-        );
-      }
+      // Kill segments - order of steps 1, 2, 3, 4 must remain the same
 
-      // Nuke Segments
-      taskActionClient.submit(new SegmentNukeAction(new 
HashSet<>(unusedSegments)));
-      emitMetric(toolbox.getEmitter(), 
TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, unusedSegments.size());
+      // 1. Determine parent segment ids of killable unused segments
+      final Map<String, String> upgradedFromSegmentIds
+          = fetchParentIdsForSegments(toolbox, unusedSegmentsPlus);
 
-      // Determine segments to be killed
-      final List<DataSegment> segmentsToBeKilled
-          = getKillableSegments(unusedSegments, upgradedFromSegmentIds, 
usedSegmentLoadSpecs, taskActionClient);
+      // 2. Identify killable segments whose load specs are not shared with 
any other segment
+      final List<DataSegment> segmentsToKillFromDeepStore = 
getKillableSegments(
+          unusedSegments,
+          upgradedFromSegmentIds,
+          usedSegmentLoadSpecs,
+          taskActionClient
+      );
 
+      // 2a. Track segments that cannot be removed from deep store yet
       final Set<DataSegment> segmentsNotKilled = new HashSet<>(unusedSegments);
-      segmentsToBeKilled.forEach(segmentsNotKilled::remove);
-
+      segmentsToKillFromDeepStore.forEach(segmentsNotKilled::remove);
       if (!segmentsNotKilled.isEmpty()) {
         LOG.warn(
-            "Skipping kill of [%d] segments from deep storage as their load 
specs are used by other segments.",
-            segmentsNotKilled.size()
+            "Skipping kill of [%d] segments of datasource[%s] from deep 
storage"
+            + " as their load specs are shared by other segments.",

Review Comment:
   I think instead of "other", it'd help to provide the summary of the upgraded 
segment set



##########
indexing-service/src/main/java/org/apache/druid/indexing/common/task/KillUnusedSegmentsTask.java:
##########
@@ -352,7 +347,47 @@ protected List<DataSegment> 
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<String, String> fetchParentIdsForSegments(
+      TaskToolbox toolbox,
+      List<DataSegmentPlus> unusedSegments
+  )
+  {
+    try {
+      final Set<String> segmentIds = unusedSegments.stream().map(
+          s -> s.getDataSegment().getId().toString()
+      ).collect(Collectors.toSet());
+
+      return toolbox.getTaskActionClient().submit(
+          new RetrieveUpgradedFromSegmentIdsAction(getDataSource(), segmentIds)
+      ).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."

Review Comment:
   What do you think about adding a unit test in `KillUnusedSegmentsTaskTest` 
or `UnusedSegmentsKillerTest` to verify that segments aren't spuriously purged 
in this case?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to