diff --git a/benchmarks/src/test/java/org/apache/druid/benchmark/indexing/SqlSegmentsMetadataQueryBenchmark.java b/benchmarks/src/test/java/org/apache/druid/benchmark/indexing/SqlSegmentsMetadataQueryBenchmark.java index aa49884cd0bf..13fcd970b36b 100644 --- a/benchmarks/src/test/java/org/apache/druid/benchmark/indexing/SqlSegmentsMetadataQueryBenchmark.java +++ b/benchmarks/src/test/java/org/apache/druid/benchmark/indexing/SqlSegmentsMetadataQueryBenchmark.java @@ -25,6 +25,7 @@ import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.metadata.IndexerSqlMetadataStorageCoordinatorTestBase; import org.apache.druid.metadata.MetadataStorageTablesConfig; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; import org.apache.druid.metadata.SqlSegmentsMetadataQuery; import org.apache.druid.metadata.TestDerbyConnector; import org.apache.druid.segment.TestDataSource; @@ -137,7 +138,13 @@ private Set readAsSet(Function { final SqlSegmentsMetadataQuery query = - SqlSegmentsMetadataQuery.forHandle(handle, derbyConnector, tablesConfig, TestHelper.JSON_MAPPER); + SqlSegmentsMetadataQuery.forHandle( + handle, + derbyConnector, + tablesConfig, + new SegmentsMetadataManagerConfig(null, null, null), + TestHelper.JSON_MAPPER + ); try (CloseableIterator iterator = iterableReader.apply(query)) { return ImmutableSet.copyOf(iterator); diff --git a/docs/configuration/index.md b/docs/configuration/index.md index d48982226649..238f0b357faa 100644 --- a/docs/configuration/index.md +++ b/docs/configuration/index.md @@ -1040,7 +1040,9 @@ None of the configs that apply to [auto-kill performed by the Coordinator](../da |Property|Description|Default| |--------|-----------|-------| |`druid.manager.segments.killUnused.enabled`|Boolean flag to enable auto-kill of eligible unused segments on the Overlord. This feature can be used only when [segment metadata caching](#segment-metadata-cache) is enabled on the Overlord and MUST NOT be enabled if `druid.coordinator.kill.on` is already set to `true` on the Coordinator.|`true`| -|`druid.manager.segments.killUnused.bufferPeriod`|Period after which a segment marked as unused becomes eligible for auto-kill on the Overlord. This config is effective only if `druid.manager.segments.killUnused.enabled` is set to `true`.|`P30D` (30 days)| +|`druid.manager.segments.killUnused.bufferPeriod`|ISO8601 Period after which a segment marked as unused cannot be marked as used anymore and becomes eligible for auto-kill on the Overlord. This config is effective only if `druid.manager.segments.killUnused.enabled` is set to `true`.|`P30D` (30 days)| +|`druid.manager.segments.killUnused.dutyPeriod`|ISO8601 Period defining the frequency at which unused segments should be added to the kill queue. If the queue already has some unused segments, new segments are not added. This config is effective only if `druid.manager.segments.killUnused.enabled` is set to `true`.|`PT1H` (1 hour)| +|`druid.manager.segments.killUnused.maxSegmentsToKill`|Maximum number of unused segments that can be added to the kill queue in a single cycle. A very large value for this config may cause the metadata store to slow down while fetching unused segments to kill, whereas a very small value would cause the kill operation to be ineffective as it wouldn't be able to catch up with the number of old unused segments in the cluster. This config is effective only if `druid.manager.segments.killUnused.enabled` is set to `true`.|200,000| #### Overlord dynamic configuration diff --git a/docs/data-management/delete.md b/docs/data-management/delete.md index 799e8b4b8b91..54658097c000 100644 --- a/docs/data-management/delete.md +++ b/docs/data-management/delete.md @@ -143,7 +143,7 @@ These embedded tasks offer several advantages over auto-kill performed by the Co - run on the Overlord and do not take up task slots. - finish faster as they save on the overhead of launching a task process. - kill a small number of segments per task, to ensure that locks on an interval are not held for too long. -- skip locked intervals to avoid head-of-line blocking in kill tasks. +- use a dedicated `KILL` lock type which allows killing of segments in the background while an ingestion proceeds for the same interval. - require little to no configuration. - can keep up with a large number of unused segments in the cluster. - take advantage of the segment metadata cache on the Overlord. diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java index bc7385e6bf6e..49be8522ca01 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java @@ -30,6 +30,7 @@ import org.apache.druid.indexing.kafka.supervisor.KafkaSupervisorSpec; import org.apache.druid.indexing.overlord.Segments; import org.apache.druid.java.util.common.HumanReadableBytes; +import org.apache.druid.metadata.UnusedSegmentKillerConfig; import org.apache.druid.query.DruidMetrics; import org.apache.druid.rpc.UpdateResponse; import org.apache.druid.rpc.indexing.OverlordClient; @@ -105,6 +106,8 @@ public void stop() } }; + final Period killBufferPeriod = Period.millis(100).minus(UnusedSegmentKillerConfig.GRACE_PERIOD); + indexer.setServerMemory(1_000_000_000L) .addProperty("druid.segment.handoff.pollDuration", "PT0.1s") .addProperty("druid.worker.capacity", "10"); @@ -112,7 +115,7 @@ public void stop() .addProperty("druid.manager.segments.useIncrementalCache", "ifSynced") .addProperty("druid.manager.segments.pollDuration", "PT0.1s") .addProperty("druid.manager.segments.killUnused.enabled", "true") - .addProperty("druid.manager.segments.killUnused.bufferPeriod", "PT0.1s") + .addProperty("druid.manager.segments.killUnused.bufferPeriod", killBufferPeriod.toString()) .addProperty("druid.manager.segments.killUnused.dutyPeriod", "PT1s"); coordinator.addProperty("druid.manager.segments.useIncrementalCache", "ifSynced"); cluster.addExtension(KafkaIndexTaskModule.class) @@ -203,7 +206,7 @@ public void test_ingest10kRows_ofSelfClusterMetrics_andVerifyValues() @MethodSource("getCompactionSupervisorTestParams") @ParameterizedTest(name = "engine={0}, policy={1}") @Timeout(120) - public void test_ingestClusterMetrics_withConcurrentCompactionSupervisor_andSkipKillOfUnusedSegments( + public void test_ingestClusterMetrics_withConcurrentCompactionSupervisor_andKillUnusedSegments( CompactionEngine engine, CompactionCandidateSearchPolicy policy ) @@ -303,9 +306,16 @@ public void test_ingestClusterMetrics_withConcurrentCompactionSupervisor_andSkip agg -> agg.hasSumAtLeast(1) ); - // Verify that the segments are skipped since the interval is still being appended to + // Verify that some unused segments have been killed from metadata store + overlord.latchableEmitter().waitForEventAggregate( + event -> event.hasMetricName("segment/killed/metadataStore/count") + .hasDimension(DruidMetrics.DATASOURCE, dataSource), + agg -> agg.hasSumAtLeast(1) + ); + + // Verify that some segments were not deleted from deep store due to being upgraded overlord.latchableEmitter().waitForEventAggregate( - event -> event.hasMetricName("segment/kill/skippedIntervals/count") + event -> event.hasMetricName("segment/kill/deepStorageSkipped/count") .hasDimension(DruidMetrics.DATASOURCE, dataSource), agg -> agg.hasSumAtLeast(1) ); diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/common/TaskLockType.java b/indexing-service/src/main/java/org/apache/druid/indexing/common/TaskLockType.java index 43f79cf26511..964e4ef62686 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/common/TaskLockType.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/common/TaskLockType.java @@ -24,6 +24,7 @@ * 2) An appending task may use only EXCLUSIVE, SHARED or APPEND locks * 3) A replacing task may use only EXCLUSIVE or REPLACE locks * 4) REPLACE and APPEND locks can only be used with timechunk locking + * 5) Only kill tasks (task type "kill") may acquire a KILL lock */ public enum TaskLockType { @@ -49,5 +50,12 @@ public enum TaskLockType * and with at most one REPLACE lock whose interval encloses that of the APPEND lock. * They are incompatible with all other active locks. */ - APPEND + APPEND, + /** + * There can be at most one active KILL lock for a given interval. + * It can co-exist with any other active lock type (EXCLUSIVE, SHARED, REPLACE, APPEND), + * but not with another KILL lock whose interval overlaps. + * Only tasks of type "kill" may acquire a KILL lock. + */ + KILL } 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 e31ac044d89b..5c74770c2305 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 @@ -315,7 +315,16 @@ public TaskStatus runTask(TaskToolbox toolbox) throws Exception // 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()); + + final int numSegmentsDeletedFromDeepStore = segmentsToKillFromDeepStore.size(); + emitMetric(toolbox.getEmitter(), TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE, numSegmentsDeletedFromDeepStore); + if (numSegmentsDeletedFromMetadataStore > numSegmentsDeletedFromDeepStore) { + emitMetric( + toolbox.getEmitter(), + TaskMetrics.SEGMENTS_SKIPPED_DEEPSTORE_KILL, + numSegmentsDeletedFromMetadataStore - numSegmentsDeletedFromDeepStore + ); + } numBatchesProcessed++; totalSegmentsDeletedFromMetadataStore += numSegmentsDeletedFromMetadataStore; diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/TaskMetrics.java b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/TaskMetrics.java index 86456d12fdc9..71250b2a1af9 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/TaskMetrics.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/TaskMetrics.java @@ -33,5 +33,5 @@ private TaskMetrics() public static final String SEGMENTS_DELETED_FROM_METADATA_STORE = "segment/killed/metadataStore/count"; public static final String SEGMENTS_DELETED_FROM_DEEPSTORE = "segment/killed/deepStorage/count"; - public static final String FILES_DELETED_FROM_DEEPSTORE = "segment/killed/deepStorageFile/count"; + public static final String SEGMENTS_SKIPPED_DEEPSTORE_KILL = "segment/kill/deepStorageSkipped/count"; } diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/GlobalTaskLockbox.java b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/GlobalTaskLockbox.java index 8b00739167a5..6544eec28e1a 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/GlobalTaskLockbox.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/GlobalTaskLockbox.java @@ -80,6 +80,11 @@ public GlobalTaskLockbox( * Syncs the current in-memory state with the {@link TaskStorage}. * This method should be called only from {@link TaskQueue#start()}. * If the sync fails, no other operation can be performed on this lockbox. + *

+ * The sync does not restore a lock, if it is associated with a dummy task ID + * that is not persisted in the task storage (e.g. embedded kill tasks). + * This is okay since the sync happens only when Overlord becomes leader, at + * which point there wouldn't be any embedded (kill) tasks running anyway. * * @return SyncResult which needs to be processed by the caller */ diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/TaskLockbox.java b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/TaskLockbox.java index b36d3e962328..2113e5028dc9 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/TaskLockbox.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/TaskLockbox.java @@ -39,6 +39,7 @@ import org.apache.druid.indexing.common.task.PendingSegmentAllocatingTask; import org.apache.druid.indexing.common.task.Task; import org.apache.druid.indexing.common.task.Tasks; +import org.apache.druid.indexing.overlord.duty.UnusedSegmentsKiller; import org.apache.druid.java.util.common.ISE; import org.apache.druid.java.util.common.Intervals; import org.apache.druid.java.util.common.Pair; @@ -571,6 +572,14 @@ private TaskLockPosse createOrFindLockPosse(LockRequest request, Task task, bool throw new ISE("Unable to grant LockPosse to inactive Task [%s]", task.getId()); } + if (request.getType() == TaskLockType.KILL + && !UnusedSegmentsKiller.TASK_TYPE_EMBEDDED_KILL.equals(task.getType())) { + throw new ISE( + "Task[%s] of type[%s] cannot acquire a KILL lock. Only tasks of type[%s] may use KILL locks.", + task.getId(), task.getType(), UnusedSegmentsKiller.TASK_TYPE_EMBEDDED_KILL + ); + } + final TaskLockPosse posseToUse; final List foundPosses = findLockPossesOverlapsInterval( request.getInterval() @@ -587,7 +596,7 @@ private TaskLockPosse createOrFindLockPosse(LockRequest request, Task task, bool final List reusablePosses = foundPosses .stream() .filter(posse -> posse.reusableFor(request)) - .collect(Collectors.toList()); + .toList(); if (reusablePosses.isEmpty()) { // case 1) this task doesn't have any lock, but others do @@ -835,8 +844,9 @@ public void revokeLock(String taskId, TaskLock lock) throw new ISE("Cannot revoke lock for inactive task[%s]", taskId); } + // Embedded kill tasks can hold locks but are not persisted in TaskStorage final Task task = taskStorage.getTask(taskId).orNull(); - if (task == null) { + if (task == null && lock.getType() != TaskLockType.KILL) { throw new ISE("Cannot revoke lock for unknown task[%s]", taskId); } @@ -963,19 +973,20 @@ public List getLockedIntervals(List lockFilterPolici } final int priority = lockFilter.getPriority(); - final boolean isReplaceLock = TaskLockType.REPLACE.name().equals( - lockFilter.getContext().getOrDefault( - Tasks.TASK_LOCK_TYPE, - Tasks.DEFAULT_TASK_LOCK_TYPE - ) + final TaskLockType lockType = QueryContexts.getAsEnum( + Tasks.TASK_LOCK_TYPE, + lockFilter.getContext().get(Tasks.TASK_LOCK_TYPE), + TaskLockType.class, + Tasks.DEFAULT_TASK_LOCK_TYPE ); - final boolean isUsingConcurrentLocks = Boolean.TRUE.equals( - lockFilter.getContext().getOrDefault( - Tasks.USE_CONCURRENT_LOCKS, - Tasks.DEFAULT_USE_CONCURRENT_LOCKS - ) + final boolean isReplaceLock = TaskLockType.REPLACE == lockType; + final boolean isUsingConcurrentLocks = QueryContexts.getAsBoolean( + Tasks.USE_CONCURRENT_LOCKS, + lockFilter.getContext().get(Tasks.USE_CONCURRENT_LOCKS), + Tasks.DEFAULT_USE_CONCURRENT_LOCKS ); final boolean ignoreAppendLocks = isUsingConcurrentLocks || isReplaceLock; + final boolean ignoreKillLocks = lockType != TaskLockType.KILL; running.forEach( (startTime, startTimeLocks) -> startTimeLocks.forEach( @@ -983,6 +994,8 @@ public List getLockedIntervals(List lockFilterPolici taskLockPosse -> { if (taskLockPosse.getTaskLock().isRevoked()) { // do nothing + } else if (ignoreKillLocks && TaskLockType.KILL.equals(taskLockPosse.getTaskLock().getType())) { + // do nothing } else if (ignoreAppendLocks && TaskLockType.APPEND.equals(taskLockPosse.getTaskLock().getType())) { // do nothing @@ -1363,7 +1376,7 @@ Optional getOnlyTaskLockPosseContainingInterval(Task task, Interv final List filteredPosses = findLockPossesContainingInterval(interval) .stream() .filter(lockPosse -> lockPosse.containsTask(task)) - .collect(Collectors.toList()); + .toList(); if (filteredPosses.isEmpty()) { throw new ISE("Cannot find any lock for task[%s] and interval[%s]", task.getId(), interval); @@ -1422,6 +1435,8 @@ private boolean canLockCoexist(List conflictPosses, LockRequest r return canSharedLockCoexist(conflictPosses); case EXCLUSIVE: return canExclusiveLockCoexist(conflictPosses); + case KILL: + return canKillLockCoexist(conflictPosses); default: throw new UOE("Unsupported lock type: " + request.getType()); } @@ -1430,7 +1445,7 @@ private boolean canLockCoexist(List conflictPosses, LockRequest r /** * Check if an APPEND lock can coexist with a given set of conflicting posses. * An APPEND lock can coexist with any number of other APPEND locks - * OR with at most one REPLACE lock over an interval which encloes this request. + * OR with at most one REPLACE lock over an interval which encloses this request. * @param conflictPosses conflicting lock posses * @param appendRequest append lock request * @return true iff append lock can coexist with all its conflicting locks @@ -1477,6 +1492,8 @@ private boolean canReplaceLockCoexist(List conflictPosses, LockRe || posse.getTaskLock().getType().equals(TaskLockType.REPLACE)) { return false; } + // REPLACE lock can coexist with an APPEND lock only if the append interval + // is fully contained within the replace interval if (posse.getTaskLock().getType().equals(TaskLockType.APPEND) && !replaceLock.getInterval().contains(posse.getTaskLock().getInterval())) { return false; @@ -1487,7 +1504,8 @@ private boolean canReplaceLockCoexist(List conflictPosses, LockRe /** * Check if a SHARED lock can coexist with a given set of conflicting posses. - * A SHARED lock can coexist with any number of other active SHARED locks + * A SHARED lock can coexist with any number of other active SHARED locks or + * a KILL lock. * @param conflictPosses conflicting lock posses * @return true iff shared lock can coexist with all its conflicting locks */ @@ -1508,7 +1526,8 @@ private boolean canSharedLockCoexist(List conflictPosses) /** * Check if an EXCLUSIVE lock can coexist with a given set of conflicting posses. - * An EXCLUSIVE lock cannot coexist with any other overlapping active locks + * An EXCLUSIVE lock cannot coexist with any other overlapping active locks, + * except a KILL lock. * @param conflictPosses conflicting lock posses * @return true iff the exclusive lock can coexist with all its conflicting locks */ @@ -1518,11 +1537,36 @@ private boolean canExclusiveLockCoexist(List conflictPosses) if (posse.getTaskLock().isRevoked()) { continue; } - return false; + if (!isLockTypeKill(posse)) { + return false; + } } return true; } + /** + * Check if a KILL lock can coexist with a given set of conflicting posses. + * A KILL lock can coexist with any other lock type but not with another KILL lock. + * @param conflictPosses conflicting lock posses + * @return true iff the kill lock can coexist with all its conflicting locks + */ + private boolean canKillLockCoexist(List conflictPosses) + { + for (TaskLockPosse posse : conflictPosses) { + if (posse.getTaskLock().isRevoked()) { + continue; + } + if (isLockTypeKill(posse)) { + return false; + } + } + return true; + } + + private boolean isLockTypeKill(TaskLockPosse posse) + { + return posse.getTaskLock().getType().equals(TaskLockType.KILL); + } /** * Verify if every incompatible active lock is revokable. If yes, revoke all of them. @@ -1545,7 +1589,9 @@ private boolean revokeAllIncompatibleActiveLocksIfPossible( final List possesToRevoke = new ArrayList<>(); for (TaskLockPosse posse : conflictPosses) { - if (posse.getTaskLock().isRevoked()) { + // No need to revoke an already revoked lock or a KILL lock (unless the new LockRequest is also a KILL) + if (posse.getTaskLock().isRevoked() + || (isLockTypeKill(posse) && type != TaskLockType.KILL)) { continue; } switch (type) { @@ -1582,6 +1628,15 @@ private boolean revokeAllIncompatibleActiveLocksIfPossible( possesToRevoke.add(posse); } break; + case KILL: + // KILL locks are incompatible only with other KILL locks + if (isLockTypeKill(posse)) { + if (posse.getTaskLock().getNonNullPriority() >= priority) { + return false; + } + possesToRevoke.add(posse); + } + break; default: throw new UOE("Unsupported lock type: " + type); } @@ -1686,6 +1741,7 @@ boolean reusableFor(LockRequest request) case REPLACE: case APPEND: case SHARED: + case KILL: if (request instanceof TimeChunkLockRequest) { return taskLock.getInterval().contains(request.getInterval()) && taskLock.getGroupId().equals(request.getGroupId()); 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 a5f8ffa9dc96..2cfa95004ad3 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 @@ -25,6 +25,7 @@ import org.apache.druid.common.utils.IdUtils; import org.apache.druid.discovery.DruidLeaderSelector; import org.apache.druid.indexer.TaskStatus; +import org.apache.druid.indexing.common.TaskLockType; import org.apache.druid.indexing.common.TaskToolbox; import org.apache.druid.indexing.common.actions.TaskActionClient; import org.apache.druid.indexing.common.actions.TaskActionClientFactory; @@ -34,7 +35,6 @@ import org.apache.druid.indexing.common.task.Tasks; import org.apache.druid.indexing.overlord.GlobalTaskLockbox; import org.apache.druid.indexing.overlord.IndexerMetadataStorageCoordinator; -import org.apache.druid.indexing.overlord.config.DefaultTaskConfig; import org.apache.druid.java.util.common.DateTimes; import org.apache.druid.java.util.common.Stopwatch; import org.apache.druid.java.util.common.concurrent.ScheduledExecutorFactory; @@ -76,6 +76,12 @@ public class UnusedSegmentsKiller implements OverlordDuty private static final String TASK_ID_PREFIX = "overlord-issued"; + /** + * Task type for embedded kill tasks. This type is not registered as a valid + * JSON subtype of {@code Task} since embedded kill tasks are never serialized. + */ + public static final String TASK_TYPE_EMBEDDED_KILL = "kill_embedded"; + private static final int INITIAL_KILL_QUEUE_SIZE = 1000; private static final int MAX_INTERVALS_TO_KILL = 10_000; private static final int MAX_SEGMENTS_TO_KILL_IN_BATCH = 1000; @@ -125,7 +131,6 @@ public class UnusedSegmentsKiller implements OverlordDuty @Inject public UnusedSegmentsKiller( SegmentsMetadataManagerConfig config, - DefaultTaskConfig defaultTaskConfig, TaskActionClientFactory taskActionClientFactory, IndexerMetadataStorageCoordinator storageCoordinator, @IndexingService DruidLeaderSelector leaderSelector, @@ -260,7 +265,7 @@ private void rebuildKillQueue() // Identify intervals with unused segments which are eligible for kill final Map eligibleIntervals = storageCoordinator.retrieveSomeUnusedSegmentIntervals( - DateTimes.nowUtc().minus(killConfig.getBufferPeriod()), + killConfig.getMaxUpdatedTimeOfKillableSegment(), MAX_INTERVALS_TO_KILL, killConfig.getMaxSegmentsToKill() ); @@ -372,7 +377,7 @@ private void runKillTask(KillCandidate candidate, String taskId) final EmbeddedKillTask killTask = new EmbeddedKillTask( taskId, candidate, - DateTimes.nowUtc().minus(killConfig.getBufferPeriod()) + killConfig.getMaxUpdatedTimeOfKillableSegment() ); final TaskActionClient taskActionClient = taskActionClientFactory.create(killTask); @@ -474,13 +479,22 @@ private EmbeddedKillTask( candidate.dataSource(), candidate.interval(), null, - Map.of(Tasks.PRIORITY_KEY, Tasks.DEFAULT_EMBEDDED_KILL_TASK_PRIORITY), + Map.of( + Tasks.PRIORITY_KEY, Tasks.DEFAULT_EMBEDDED_KILL_TASK_PRIORITY, + Tasks.TASK_LOCK_TYPE, TaskLockType.KILL.name() + ), MAX_SEGMENTS_TO_KILL_IN_BATCH, candidate.numSegmentsToKill(), maxUpdatedTimeOfEligibleSegment ); } + @Override + public String getType() + { + return TASK_TYPE_EMBEDDED_KILL; + } + @Override protected List fetchNextBatchOfUnusedSegments(TaskToolbox toolbox, int nextBatchSize) { 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 8721148083c3..3fc1d9d188a6 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 @@ -302,9 +302,11 @@ private SqlSegmentMetadataTransactionFactory setupTransactionFactory( ? SegmentMetadataCache.UsageMode.ALWAYS : SegmentMetadataCache.UsageMode.NEVER; + final SegmentsMetadataManagerConfig managerConfig = + new SegmentsMetadataManagerConfig(Period.seconds(1), cacheMode, null); segmentMetadataCache = new HeapMemorySegmentMetadataCache( objectMapper, - Suppliers.ofInstance(new SegmentsMetadataManagerConfig(Period.seconds(1), cacheMode, null)), + Suppliers.ofInstance(managerConfig), Suppliers.ofInstance(metadataStorageTablesConfig), new NoopSegmentSchemaCache(), new IndexingStateCache(), @@ -322,6 +324,7 @@ private SqlSegmentMetadataTransactionFactory setupTransactionFactory( testDerbyConnector, leaderSelector, segmentMetadataCache, + managerConfig, emitter ) { 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 090ed20269a5..0c66a619faa0 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 @@ -336,9 +336,11 @@ private SqlSegmentMetadataTransactionFactory createTransactionFactory() = useSegmentMetadataCache ? SegmentMetadataCache.UsageMode.ALWAYS : SegmentMetadataCache.UsageMode.NEVER; + final SegmentsMetadataManagerConfig managerConfig = + new SegmentsMetadataManagerConfig(Period.millis(10), cacheMode, null); segmentMetadataCache = new HeapMemorySegmentMetadataCache( objectMapper, - Suppliers.ofInstance(new SegmentsMetadataManagerConfig(Period.millis(10), cacheMode, null)), + Suppliers.ofInstance(managerConfig), derbyConnectorRule.metadataTablesConfigSupplier(), segmentSchemaCache, indexingStateCache, @@ -356,6 +358,7 @@ private SqlSegmentMetadataTransactionFactory createTransactionFactory() derbyConnectorRule.getConnector(), leaderSelector, segmentMetadataCache, + managerConfig, NoopServiceEmitter.instance() ); } diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/CommandQueueTask.java b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/CommandQueueTask.java index 5945d1488c38..a9e5607c89de 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/CommandQueueTask.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/CommandQueueTask.java @@ -19,6 +19,8 @@ package org.apache.druid.indexing.common.task.concurrent; +import com.google.common.base.Throwables; +import org.apache.druid.error.DruidException; import org.apache.druid.indexer.TaskStatus; import org.apache.druid.indexing.common.TaskToolbox; import org.apache.druid.indexing.common.actions.TaskActionClient; @@ -137,7 +139,12 @@ private V waitForCommandToFinish(Command command) return command.value.get(10, TimeUnit.SECONDS); } catch (Exception e) { - throw new ISE(e, "Error waiting for command on task[%s] to finish", getId()); + final Throwable rootCause = Throwables.getRootCause(e); + if (rootCause instanceof DruidException druidException) { + throw druidException; + } else { + throw new ISE(e, "Error waiting for command on task[%s] to finish", getId()); + } } } diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndAppendTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndAppendTest.java index 2aa0895c9618..6a11cc30ce4d 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndAppendTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndAppendTest.java @@ -19,7 +19,6 @@ package org.apache.druid.indexing.common.task.concurrent; -import com.google.common.base.Throwables; import com.google.common.collect.ImmutableList; import com.google.common.collect.Iterables; import com.google.common.collect.Sets; @@ -590,14 +589,13 @@ public void testAllocateLockReplaceDayAppendMonth() // Verify that segment cannot be committed since there is no lock final DataSegment segmentV10 = createSegment(FIRST_OF_JAN_23, SEGMENT_V0); - final ISE exception = Assertions.assertThrows(ISE.class, () -> replaceTask.commitReplaceSegments(segmentV10)); - final Throwable throwable = Throwables.getRootCause(exception); + final DruidException exception = Assertions.assertThrows(DruidException.class, () -> replaceTask.commitReplaceSegments(segmentV10)); Assertions.assertEquals( StringUtils.format( "Segment IDs[[%s]] are not covered by locks[[]] for task[%s]", segmentV10.getId(), replaceTask.getId() ), - throwable.getMessage() + exception.getMessage() ); final DataSegment segmentV01 = asSegment(pendingSegment); @@ -658,9 +656,9 @@ public void testLockAllocateAppendYearReplaceQuarter() final DataSegment segmentV1Q3 = createSegment(JUL_AUG_SEP_23, replaceLock.getVersion()); final DataSegment segmentV1Q4 = createSegment(OCT_NOV_DEC_23, replaceLock.getVersion()); - Assertions.assertFalse( - replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) - .isSuccess() + Assertions.assertThrows( + DruidException.class, + () -> replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) ); verifyIntervalHasUsedSegments(YEAR_23, segmentV01); @@ -683,9 +681,9 @@ public void testLockAllocateReplaceQuarterAppendYear() final DataSegment segmentV1Q3 = createSegment(JUL_AUG_SEP_23, replaceLock.getVersion()); final DataSegment segmentV1Q4 = createSegment(OCT_NOV_DEC_23, replaceLock.getVersion()); - Assertions.assertFalse( - replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) - .isSuccess() + Assertions.assertThrows( + DruidException.class, + () -> replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) ); final DataSegment segmentV01 = asSegment(pendingSegment); @@ -711,9 +709,9 @@ public void testAllocateLockReplaceQuarterAppendYear() final DataSegment segmentV1Q3 = createSegment(JUL_AUG_SEP_23, replaceLock.getVersion()); final DataSegment segmentV1Q4 = createSegment(OCT_NOV_DEC_23, replaceLock.getVersion()); - Assertions.assertFalse( - replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) - .isSuccess() + Assertions.assertThrows( + DruidException.class, + () -> replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) ); final DataSegment segmentV01 = asSegment(pendingSegment); @@ -745,9 +743,9 @@ public void testAllocateLockAppendYearReplaceQuarter() final DataSegment segmentV1Q3 = createSegment(JUL_AUG_SEP_23, replaceLock.getVersion()); final DataSegment segmentV1Q4 = createSegment(OCT_NOV_DEC_23, replaceLock.getVersion()); - Assertions.assertFalse( - replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) - .isSuccess() + Assertions.assertThrows( + DruidException.class, + () -> replaceTask.commitReplaceSegments(segmentV1Q1, segmentV1Q2, segmentV1Q3, segmentV1Q4) ); verifyIntervalHasUsedSegments(YEAR_23, segmentV01); @@ -1203,13 +1201,9 @@ public void test_allocateCommitDelete_createsFreshVersion_uptoMaxAllowedRetries( } // Verify that the next attempt fails - final ISE exception = Assertions.assertThrows( - ISE.class, - () -> appendTask.allocateSegmentForTimestamp(FIRST_OF_JAN_23.getStart(), Granularities.DAY) - ); - final DruidException rootCause = Assertions.assertInstanceOf( + final DruidException rootCause = Assertions.assertThrows( DruidException.class, - Throwables.getRootCause(exception) + () -> appendTask.allocateSegmentForTimestamp(FIRST_OF_JAN_23.getStart(), Granularities.DAY) ); Assertions.assertEquals(DruidException.Persona.OPERATOR, rootCause.getTargetPersona()); Assertions.assertEquals(DruidException.Category.RUNTIME_FAILURE, rootCause.getCategory()); diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndStreamingAppendTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndStreamingAppendTest.java index 7bc0b6685617..57c25a650603 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndStreamingAppendTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/concurrent/ConcurrentReplaceAndStreamingAppendTest.java @@ -23,6 +23,7 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.Iterables; import com.google.common.collect.Sets; +import org.apache.druid.error.DruidException; import org.apache.druid.indexing.common.MultipleFileTaskReportFileWriter; import org.apache.druid.indexing.common.TaskLock; import org.apache.druid.indexing.common.TaskStorageDirTracker; @@ -96,7 +97,7 @@ *

  • LOCK: Acquisition of a lock on an interval by a replace task
  • *
  • ALLOCATE: Allocation of a pending segment by an append task
  • *
  • REPLACE: Commit of segments created by a replace task
  • - *
  • APPEND: Commit of segments created by an append task
  • + *
  • APPEND: Commit of segments created by an append task
  • * */ public class ConcurrentReplaceAndStreamingAppendTest extends IngestionTestBase @@ -606,7 +607,7 @@ public void testAllocateLockReplaceDayAppendMonth() // Verify that segment cannot be committed since there is no lock final DataSegment segmentV10 = createSegment(FIRST_OF_JAN_23, SEGMENT_V0); - final ISE exception = Assertions.assertThrows(ISE.class, () -> commitReplaceSegments(segmentV10)); + final DruidException exception = Assertions.assertThrows(DruidException.class, () -> commitReplaceSegments(segmentV10)); final Throwable throwable = Throwables.getRootCause(exception); Assertions.assertEquals( StringUtils.format( diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/GlobalTaskLockboxTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/GlobalTaskLockboxTest.java index 42e6e365b306..c2f9f25e6890 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/GlobalTaskLockboxTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/GlobalTaskLockboxTest.java @@ -29,6 +29,7 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import com.google.common.collect.Iterables; +import org.apache.druid.common.utils.IdUtils; import org.apache.druid.error.ExceptionMatcher; import org.apache.druid.indexer.TaskStatus; import org.apache.druid.indexing.common.LockGranularity; @@ -44,6 +45,7 @@ import org.apache.druid.indexing.common.task.NoopTask; import org.apache.druid.indexing.common.task.Task; import org.apache.druid.indexing.common.task.Tasks; +import org.apache.druid.indexing.overlord.duty.UnusedSegmentsKiller; import org.apache.druid.jackson.DefaultObjectMapper; import org.apache.druid.java.util.common.DateTimes; import org.apache.druid.java.util.common.ISE; @@ -55,9 +57,11 @@ import org.apache.druid.metadata.IndexerSQLMetadataStorageCoordinator; import org.apache.druid.metadata.LockFilterPolicy; import org.apache.druid.metadata.MetadataStorageTablesConfig; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; import org.apache.druid.metadata.TestDerbyConnector; import org.apache.druid.metadata.segment.SqlSegmentMetadataTransactionFactory; import org.apache.druid.metadata.segment.cache.NoopSegmentMetadataCache; +import org.apache.druid.segment.TestDataSource; import org.apache.druid.segment.TestHelper; import org.apache.druid.segment.metadata.CentralizedDatasourceSchemaConfig; import org.apache.druid.segment.metadata.HeapMemoryIndexingStateStorage; @@ -140,6 +144,7 @@ public void setup() derbyConnector, new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ), objectMapper, @@ -487,6 +492,7 @@ public void testSyncWithUnknownTaskTypesFromModuleNotLoaded() derbyConnector, new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ), loadedMapper, @@ -1308,6 +1314,31 @@ public void testGetLockedIntervalsForLowerPriorityUseConcurrentLocks() Assert.assertTrue(conflictingIntervals.isEmpty()); } + @Test + public void test_getLockedIntervals_withOngoingKill_returnsEmpty() + { + final Interval killInterval = Intervals.of("2017/2018"); + final Task task = newEmbeddedKillTask(HIGH_PRIORITY); + lockbox.add(task); + taskStorage.insert(task, TaskStatus.running(task.getId())); + tryTimeChunkLock( + TaskLockType.KILL, + task, + killInterval + ); + + LockFilterPolicy requestForExclusiveLowerPriorityLock = new LockFilterPolicy( + task.getDataSource(), + 25, + null, + null + ); + + Map> conflictingIntervals = + lockbox.getLockedIntervals(ImmutableList.of(requestForExclusiveLowerPriorityLock)); + Assert.assertTrue(conflictingIntervals.isEmpty()); + } + @Test public void testGetActiveLocks() @@ -1848,6 +1879,199 @@ public void testReplaceLockCanRevokeAllIncompatible() validator.expectRevokedLocks(appendLock0, appendLock2, exclusiveLock, replaceLock, sharedLock); } + @Test + public void testKillLockCompatibility() + { + final Task killTask = newEmbeddedKillTask(MEDIUM_PRIORITY); + final TaskLock theLock = validator.expectKillLockCreated(killTask, Intervals.of("2017/2018")); + + // A KILL lock cannot coexist with another KILL lock on an overlapping interval + validator.expectKillLockNotGranted(newEmbeddedKillTask(MEDIUM_PRIORITY), Intervals.of("2017-05-01/2017-06-01")); + + // A KILL lock can coexist with all other lock types + final TaskLock exclusiveLock = validator.expectLockCreated( + TaskLockType.EXCLUSIVE, + Intervals.of("2017-05-01/2017-06-01"), + MEDIUM_PRIORITY + ); + validator.expectActiveLocks(theLock, exclusiveLock); + validator.expectRevokedLocks(); + + final TaskLock sharedLock = validator.expectLockCreated( + TaskLockType.SHARED, + Intervals.of("2017-05-01/2017-06-01"), + MEDIUM_PRIORITY + 1 + ); + validator.expectActiveLocks(theLock, sharedLock); + validator.expectRevokedLocks(exclusiveLock); + + final TaskLock replaceLock = validator.expectLockCreated( + TaskLockType.REPLACE, + Intervals.of("2017-03-01/2017-09-01"), + MEDIUM_PRIORITY + 2 + ); + final TaskLock appendLock = validator.expectLockCreated( + TaskLockType.APPEND, + Intervals.of("2017-05-01/2017-06-01"), + MEDIUM_PRIORITY + 3 + ); + validator.expectActiveLocks(theLock, replaceLock, appendLock); + validator.expectRevokedLocks(exclusiveLock, sharedLock); + } + + @Test + public void testKillLockCanRevokeIncompatibleKillLock() + { + final TaskLock lowPriorityKillLock = validator.expectKillLockCreated( + newEmbeddedKillTask(LOW_PRIORITY), + Intervals.of("2017-05-01/2017-06-01") + ); + + // A higher-priority KILL lock can revoke a lower-priority KILL lock + final TaskLock highPriorityKillLock = validator.expectKillLockCreated( + newEmbeddedKillTask(HIGH_PRIORITY), + Intervals.of("2017/2018") + ); + + validator.expectActiveLocks(highPriorityKillLock); + validator.expectRevokedLocks(lowPriorityKillLock); + } + + @Test + public void testKillLockCannotRevokeHigherPriorityKillLock() + { + validator.expectKillLockCreated(newEmbeddedKillTask(HIGH_PRIORITY), Intervals.of("2017-05-01/2017-06-01")); + validator.expectKillLockNotGranted(newEmbeddedKillTask(LOW_PRIORITY), Intervals.of("2017/2018")); + } + + @Test + public void testOnlyEmbeddedKillTaskCanAcquireKillLock() + { + final Task nonKillTask = NoopTask.ofPriority(MEDIUM_PRIORITY); + lockbox.add(nonKillTask); + taskStorage.insert(nonKillTask, TaskStatus.running(nonKillTask.getId())); + + Assert.assertThrows( + ISE.class, + () -> lockbox.tryLock( + nonKillTask, + new TimeChunkLockRequest(TaskLockType.KILL, nonKillTask, Intervals.of("2017/2018"), null) + ) + ); + } + + @Test + public void testKillLockAcquireAndReleaseWithoutTaskInStorage() + { + // Acquire a KILL lock on an interval + final Task killTask = newEmbeddedKillTask(MEDIUM_PRIORITY); + final Interval lockInterval = Intervals.of("2017/2018"); + validator.expectKillLockCreated(killTask, lockInterval); + + // Verify that the lock blocks another KILL on an overlapping interval + validator.expectKillLockNotGranted(newEmbeddedKillTask(MEDIUM_PRIORITY), lockInterval); + + // Release the original task and its lock so that other tasks can acquire a lock + lockbox.remove(killTask); + validator.expectKillLockCreated(newEmbeddedKillTask(MEDIUM_PRIORITY), lockInterval); + } + + @Test + public void testKillLockNotRestoredAfterSyncFromStorage() + { + // Acquire a KILL lock for a task that was never inserted into taskStorage + final Interval lockInterval = Intervals.of("2017/2018"); + final Task killTask = newEmbeddedKillTask(MEDIUM_PRIORITY); + final TaskLock killLock = validator.expectKillLockCreated(killTask, lockInterval); + + // Verify that other tasks cannot acquire a lock on the same interval + validator.expectKillLockNotGranted(newEmbeddedKillTask(MEDIUM_PRIORITY), lockInterval); + + // Sync from storage — the kill task was not persisted, so the lock must disappear + final GlobalTaskLockbox newBox = new GlobalTaskLockbox(taskStorage, metadataStorageCoordinator); + final TaskLockboxSyncResult result = newBox.syncFromStorage(); + Assert.assertEquals(0, result.getTaskLockCount()); + + // Verify that the lock is still present in storage + final List locksInStorage = taskStorage.getLocks(killTask.getId()); + Assert.assertEquals(1, locksInStorage.size()); + Assert.assertEquals(killLock, locksInStorage.getFirst()); + + // A new KILL lock on the same interval should now be grantable + this.lockbox = newBox; + new TaskLockboxValidator(newBox, taskStorage) + .expectKillLockCreated(newEmbeddedKillTask(MEDIUM_PRIORITY), lockInterval); + } + + @Test + public void testKillLockRevocationByHigherPriorityKillTaskNotInStorage() + { + // Low-priority kill task not in storage acquires a KILL lock + final Task lowPriorityKillTask = newEmbeddedKillTask(LOW_PRIORITY); + lockbox.add(lowPriorityKillTask); + final LockResult lowResult = lockbox.tryLock( + lowPriorityKillTask, + new TimeChunkLockRequest(TaskLockType.KILL, lowPriorityKillTask, Intervals.of("2017-06-01/2017-07-01"), null) + ); + Assert.assertTrue(lowResult.isOk()); + Assert.assertFalse(lowResult.getTaskLock().isRevoked()); + + // High-priority kill task not in storage revokes the low-priority KILL lock + final Task highPriorityKillTask = newEmbeddedKillTask(HIGH_PRIORITY); + lockbox.add(highPriorityKillTask); + final LockResult highResult = lockbox.tryLock( + highPriorityKillTask, + new TimeChunkLockRequest(TaskLockType.KILL, highPriorityKillTask, Intervals.of("2017/2018"), null) + ); + Assert.assertTrue(highResult.isOk()); + Assert.assertFalse(highResult.getTaskLock().isRevoked()); + + // Re-acquiring the low-priority lock returns the revoked copy + final LockResult revokedResult = lockbox.tryLock( + lowPriorityKillTask, + new TimeChunkLockRequest(TaskLockType.KILL, lowPriorityKillTask, Intervals.of("2017-06-01/2017-07-01"), null) + ); + Assert.assertFalse(revokedResult.isOk()); + Assert.assertTrue(revokedResult.getTaskLock().isRevoked()); + + lockbox.remove(lowPriorityKillTask); + lockbox.remove(highPriorityKillTask); + } + + @Test + public void testKillLockCoexistsWithOtherLocksNotInStorage() + { + // Acquire an EXCLUSIVE lock for a persisted task + final Task exclusiveTask = NoopTask.ofPriority(MEDIUM_PRIORITY); + lockbox.add(exclusiveTask); + taskStorage.insert(exclusiveTask, TaskStatus.running(exclusiveTask.getId())); + final LockResult exclusiveResult = lockbox.tryLock( + exclusiveTask, + new TimeChunkLockRequest(TaskLockType.EXCLUSIVE, exclusiveTask, Intervals.of("2017-05-01/2017-06-01"), null) + ); + Assert.assertTrue(exclusiveResult.isOk()); + + // A KILL task (not in storage) should be able to acquire a KILL lock on an overlapping interval + final Task killTask = newEmbeddedKillTask(MEDIUM_PRIORITY); + lockbox.add(killTask); + final LockResult killResult = lockbox.tryLock( + killTask, + new TimeChunkLockRequest(TaskLockType.KILL, killTask, Intervals.of("2017/2018"), null) + ); + Assert.assertTrue(killResult.isOk()); + Assert.assertFalse(killResult.getTaskLock().isRevoked()); + + // Release the kill task (not persisted) + lockbox.remove(killTask); + + // The exclusive lock should still be active + final List exclusiveLocks = taskStorage.getLocks(exclusiveTask.getId()); + Assert.assertEquals(1, exclusiveLocks.size()); + Assert.assertFalse(exclusiveLocks.get(0).isRevoked()); + + lockbox.remove(exclusiveTask); + } + @Test public void testTimechunkLockTypeTransitionForSameTaskGroup() { @@ -2070,6 +2294,25 @@ public void test_add_throwsException_ifSyncIsNotComplete() ); } + private static Task newEmbeddedKillTask(int priority) + { + return new NoopTask( + IdUtils.getRandomId(), + null, + TestDataSource.WIKI, + 1L, + 1L, + Map.of(Tasks.PRIORITY_KEY, priority) + ) + { + @Override + public String getType() + { + return UnusedSegmentsKiller.TASK_TYPE_EMBEDDED_KILL; + } + }; + } + private class TaskLockboxValidator { @@ -2123,7 +2366,7 @@ public void expectRevokedLocks(TaskLock... locks) { final Set allLocks = getAllLocks(); final Set activeLocks = getAllActiveLocks(); - Assert.assertEquals(allLocks.size() - activeLocks.size(), locks.length); + Assert.assertEquals(locks.length, allLocks.size() - activeLocks.size()); for (TaskLock lock : locks) { Assert.assertTrue(allLocks.contains(lock.revokedCopy())); Assert.assertFalse(activeLocks.contains(lock)); @@ -2134,7 +2377,7 @@ public void expectActiveLocks(TaskLock... locks) { final Set allLocks = getAllLocks(); final Set activeLocks = getAllActiveLocks(); - Assert.assertEquals(activeLocks.size(), locks.length); + Assert.assertEquals(locks.length, activeLocks.size()); for (TaskLock lock : locks) { Assert.assertTrue(allLocks.contains(lock)); Assert.assertTrue(activeLocks.contains(lock)); @@ -2145,7 +2388,9 @@ private TaskLock tryTaskLock(TaskLockType type, Task task, Interval interval) { if (tasks.add(task)) { lockbox.add(task); - taskStorage.insert(task, TaskStatus.running(task.getId())); + if (!UnusedSegmentsKiller.TASK_TYPE_EMBEDDED_KILL.equals(task.getType())) { + taskStorage.insert(task, TaskStatus.running(task.getId())); + } } TaskLock lock = tryTimeChunkLock(type, task, interval).getTaskLock(); if (lock != null) { @@ -2159,6 +2404,24 @@ private TaskLock tryTaskLock(TaskLockType type, Interval interval, int priority) return tryTaskLock(type, NoopTask.ofPriority(priority), interval); } + /** + * Adds the given {@code killTask} to the lockbox (but not to TaskStorage) + * and verifies that a kill lock is successfully granted on the specified interval. + */ + public TaskLock expectKillLockCreated(Task killTask, Interval interval) + { + final TaskLock lock = tryTaskLock(TaskLockType.KILL, killTask, interval); + Assert.assertNotNull(lock); + Assert.assertFalse(lock.isRevoked()); + return lock; + } + + public void expectKillLockNotGranted(Task killTask, Interval interval) + { + final TaskLock lock = tryTaskLock(TaskLockType.KILL, killTask, interval); + Assert.assertNull(lock); + } + private Set getAllActiveLocks() { return tasks.stream() diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLockBoxConcurrencyTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLockBoxConcurrencyTest.java index 3515790e4a2b..2ddd004e583d 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLockBoxConcurrencyTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskLockBoxConcurrencyTest.java @@ -37,6 +37,7 @@ import org.apache.druid.java.util.common.granularity.Granularities; import org.apache.druid.metadata.DerbyMetadataStorageActionHandlerFactory; import org.apache.druid.metadata.IndexerSQLMetadataStorageCoordinator; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; import org.apache.druid.metadata.TestDerbyConnector; import org.apache.druid.metadata.segment.SqlSegmentMetadataTransactionFactory; import org.apache.druid.metadata.segment.cache.NoopSegmentMetadataCache; @@ -101,6 +102,7 @@ public void setup() derbyConnector, new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ), objectMapper, diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueScaleTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueScaleTest.java index 190d78bb33c9..5432326fd3ff 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueScaleTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueScaleTest.java @@ -45,6 +45,7 @@ import org.apache.druid.java.util.common.logger.Logger; import org.apache.druid.java.util.emitter.EmittingLogger; import org.apache.druid.metadata.IndexerSQLMetadataStorageCoordinator; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; import org.apache.druid.metadata.TaskLookup; import org.apache.druid.metadata.TestDerbyConnector; import org.apache.druid.metadata.segment.SqlSegmentMetadataTransactionFactory; @@ -113,6 +114,7 @@ public void setUp() derbyConnectorRule.getConnector(), new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ), jsonMapper, 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 1c8d885c9240..bed08b4a85df 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 @@ -28,7 +28,6 @@ import org.apache.druid.indexing.overlord.GlobalTaskLockbox; import org.apache.druid.indexing.overlord.IndexerMetadataStorageCoordinator; import org.apache.druid.indexing.overlord.TimeChunkLockRequest; -import org.apache.druid.indexing.overlord.config.DefaultTaskConfig; import org.apache.druid.indexing.test.TestDataSegmentKiller; import org.apache.druid.java.util.common.ISE; import org.apache.druid.java.util.common.Intervals; @@ -83,7 +82,7 @@ public void setup() emitter = taskActionTestKit.getServiceEmitter(); leaderSelector = new TestDruidLeaderSelector(); dataSegmentKiller = new TestDataSegmentKiller(); - killerConfig = new UnusedSegmentKillerConfig(true, Period.ZERO, null, null); + killerConfig = new UnusedSegmentKillerConfig(true, zeroBufferPeriod(), null, null); killExecutor = new BlockingExecutorService("UnusedSegmentsKillerTest-%s"); storageCoordinator = taskActionTestKit.getMetadataStorageCoordinator(); initKiller(); @@ -97,7 +96,6 @@ private void initKiller() SegmentMetadataCache.UsageMode.ALWAYS, killerConfig ), - new DefaultTaskConfig(), taskActionTestKit::createTaskActionClient, storageCoordinator, leaderSelector, @@ -213,7 +211,7 @@ public void test_run_launchesEmbeddedKillTasks_ifLeader() @Test public void test_maxSegmentsKilledInRun_isLimitedByConfig() { - killerConfig = new UnusedSegmentKillerConfig(true, Period.ZERO, null, 700); + killerConfig = new UnusedSegmentKillerConfig(true, zeroBufferPeriod(), null, 700); initKiller(); leaderSelector.becomeLeader(); @@ -474,7 +472,7 @@ public void test_run_killsFutureSegment() } @Test - public void test_run_skipsLockedIntervals() throws InterruptedException + public void test_run_doesNotSkipLockedIntervals() throws InterruptedException { storageCoordinator.commitSegments(Set.copyOf(WIKI_SEGMENTS_1X10D), null); storageCoordinator.markAllSegmentsAsUnused(TestDataSource.WIKI); @@ -499,10 +497,10 @@ public void test_run_skipsLockedIntervals() throws InterruptedException killer.run(); finishQueuedKillJobs(); - // Verify that unused segments from locked intervals are not killed - emitter.verifySum(TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, 5L); - emitter.verifySum(TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE, 5L); - emitter.verifySum(UnusedSegmentsKiller.Metric.SKIPPED_INTERVALS, 5L); + // Verify that unused segments from locked intervals are also killed + emitter.verifySum(TaskMetrics.SEGMENTS_DELETED_FROM_METADATA_STORE, 10L); + emitter.verifySum(TaskMetrics.SEGMENTS_DELETED_FROM_DEEPSTORE, 10L); + emitter.verifyNotEmitted(UnusedSegmentsKiller.Metric.SKIPPED_INTERVALS); } finally { taskLockbox.remove(ingestionTask); @@ -524,4 +522,13 @@ private List retrieveUnusedSegments(Interval interval) null ); } + + /** + * Buffer period which ensures that segments are killed as soon as they become unused. + */ + private static Period zeroBufferPeriod() + { + // Subtract the grace period + return Period.ZERO.minus(UnusedSegmentKillerConfig.GRACE_PERIOD); + } } diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java index 366801ffd7ce..8a75d17a80e9 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTaskTestBase.java @@ -86,6 +86,7 @@ import org.apache.druid.java.util.metrics.StubServiceEmitter; import org.apache.druid.metadata.DerbyMetadataStorageActionHandlerFactory; import org.apache.druid.metadata.IndexerSQLMetadataStorageCoordinator; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; import org.apache.druid.metadata.TestDerbyConnector; import org.apache.druid.metadata.segment.SqlSegmentMetadataTransactionFactory; import org.apache.druid.metadata.segment.cache.NoopSegmentMetadataCache; @@ -592,6 +593,7 @@ protected void makeToolboxFactory(TestUtils testUtils, ServiceEmitter emitter, b derbyConnector, new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ), objectMapper, 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 49e18862c059..ba739e44005c 100644 --- a/server/src/main/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinator.java +++ b/server/src/main/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinator.java @@ -23,7 +23,6 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; -import com.google.common.base.Throwables; import com.google.common.collect.FluentIterable; import com.google.common.collect.ImmutableList; import com.google.common.collect.Iterables; @@ -252,9 +251,8 @@ public List retrieveUnusedSegmentsForInterval( @Nullable DateTime maxUsedStatusLastUpdatedTime ) { - final List matchingSegments = inReadOnlyDatasourceTransaction( - dataSource, - transaction -> transaction.noCacheSql().findUnusedSegments( + final List matchingSegments = inReadOnlyTransaction( + sql -> sql.findUnusedSegments( dataSource, interval, versions, @@ -2869,50 +2867,18 @@ private T inReadOnlyDatasourceTransaction( } /** - * Performs a read-only transaction using the {@link SqlSegmentsMetadataQuery}, - * which queries the metadata store directly. + * @see SegmentMetadataTransactionFactory#inReadOnlyNoCacheTransaction(Function) */ private T inReadOnlyTransaction(Function sqlQuery) { - try { - return connector.retryReadOnlyTransaction( - (handle, status) -> sqlQuery.apply( - SqlSegmentsMetadataQuery.forHandle(handle, connector, dbTables, jsonMapper) - ), - 2, 3 - ); - } - catch (Throwable t) { - Throwable rootCause = Throwables.getRootCause(t); - if (rootCause instanceof DruidException) { - throw (DruidException) rootCause; - } else { - throw t; - } - } + return transactionFactory.inReadOnlyNoCacheTransaction(sqlQuery); } /** - * Performs a write transaction using the {@link SqlSegmentsMetadataQuery}, - * which updates the metadata store directly. + * @see SegmentMetadataTransactionFactory#inReadWriteNoCacheTransaction(Function) */ private T inWriteTransaction(Function sqlUpdate) { - try { - return connector.retryTransaction( - (handle, status) -> sqlUpdate.apply( - SqlSegmentsMetadataQuery.forHandle(handle, connector, dbTables, jsonMapper) - ), - 2, 3 - ); - } - catch (Throwable t) { - Throwable rootCause = Throwables.getRootCause(t); - if (rootCause instanceof DruidException) { - throw (DruidException) rootCause; - } else { - throw t; - } - } + return transactionFactory.inReadWriteNoCacheTransaction(sqlUpdate); } } diff --git a/server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java b/server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java index 2be87f57b013..1ca0dfbe9489 100644 --- a/server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java +++ b/server/src/main/java/org/apache/druid/metadata/SQLMetadataConnector.java @@ -28,6 +28,7 @@ import org.apache.commons.codec.digest.DigestUtils; import org.apache.commons.dbcp2.BasicDataSource; import org.apache.commons.dbcp2.BasicDataSourceFactory; +import org.apache.druid.error.DruidException; import org.apache.druid.java.util.common.ISE; import org.apache.druid.java.util.common.RetryUtils; import org.apache.druid.java.util.common.StringUtils; @@ -170,7 +171,7 @@ public T retryWithHandle( ); } catch (Exception e) { - Throwables.propagateIfPossible(e); + throwIfUnchecked(e); throw new RuntimeException(e); } } @@ -193,7 +194,7 @@ public T retryTransaction(final TransactionCallback callback, final int q ); } catch (Exception e) { - Throwables.propagateIfPossible(e); + throwIfUnchecked(e); throw new RuntimeException(e); } } @@ -974,7 +975,7 @@ public final T retryReadOnlyTransaction( ); } catch (Exception e) { - Throwables.throwIfUnchecked(e); + throwIfUnchecked(e); throw new RuntimeException(e); } } @@ -1389,6 +1390,15 @@ private void validateSegmentsTable() } } + private static void throwIfUnchecked(Throwable t) + { + final Throwable rootCause = Throwables.getRootCause(t); + if (rootCause instanceof DruidException druidException) { + throw druidException; + } + Throwables.throwIfUnchecked(t); + } + public static boolean isStatementException(Throwable e) { return e instanceof StatementException || diff --git a/server/src/main/java/org/apache/druid/metadata/SegmentsMetadataManagerConfig.java b/server/src/main/java/org/apache/druid/metadata/SegmentsMetadataManagerConfig.java index d02bd6bdec28..a2e9e1b4933b 100644 --- a/server/src/main/java/org/apache/druid/metadata/SegmentsMetadataManagerConfig.java +++ b/server/src/main/java/org/apache/druid/metadata/SegmentsMetadataManagerConfig.java @@ -26,6 +26,8 @@ import org.apache.druid.metadata.segment.cache.SegmentMetadataCache; import org.joda.time.Period; +import javax.annotation.Nullable; + /** * Config that dictates polling and caching of segment metadata on leader * Coordinator or Overlord services. @@ -45,9 +47,9 @@ public class SegmentsMetadataManagerConfig @JsonCreator public SegmentsMetadataManagerConfig( - @JsonProperty("pollDuration") Period pollDuration, - @JsonProperty("useIncrementalCache") SegmentMetadataCache.UsageMode useIncrementalCache, - @JsonProperty("killUnused") UnusedSegmentKillerConfig killUnused + @JsonProperty("pollDuration") @Nullable Period pollDuration, + @JsonProperty("useIncrementalCache") @Nullable SegmentMetadataCache.UsageMode useIncrementalCache, + @JsonProperty("killUnused") @Nullable UnusedSegmentKillerConfig killUnused ) { this.pollDuration = Configs.valueOrDefault(pollDuration, Period.minutes(1)); 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 d16ba18d0b2d..d3f69baa0a01 100644 --- a/server/src/main/java/org/apache/druid/metadata/SqlSegmentsMetadataQuery.java +++ b/server/src/main/java/org/apache/druid/metadata/SqlSegmentsMetadataQuery.java @@ -55,6 +55,7 @@ import org.apache.druid.utils.CloseableUtils; import org.joda.time.DateTime; import org.joda.time.Interval; +import org.joda.time.Period; import org.skife.jdbi.v2.Handle; import org.skife.jdbi.v2.PreparedBatch; import org.skife.jdbi.v2.Query; @@ -76,12 +77,13 @@ import java.util.NoSuchElementException; import java.util.Objects; import java.util.Set; +import java.util.function.Function; import java.util.stream.Collectors; /** * An object that is used to query the segments table in the metadata store. - * Each instance of this class is scoped to a single {@link Handle} and is meant - * to be short-lived. + * Each instance of this class is scoped to a single {@link Handle} and must + * be used within a single transaction only. */ public class SqlSegmentsMetadataQuery { @@ -97,18 +99,21 @@ public class SqlSegmentsMetadataQuery private final Handle handle; private final SQLMetadataConnector connector; private final MetadataStorageTablesConfig dbTables; + private final SegmentsMetadataManagerConfig managerConfig; private final ObjectMapper jsonMapper; private SqlSegmentsMetadataQuery( final Handle handle, final SQLMetadataConnector connector, final MetadataStorageTablesConfig dbTables, + final SegmentsMetadataManagerConfig managerConfig, final ObjectMapper jsonMapper ) { this.handle = handle; this.connector = connector; this.dbTables = dbTables; + this.managerConfig = managerConfig; this.jsonMapper = jsonMapper; } @@ -120,10 +125,11 @@ public static SqlSegmentsMetadataQuery forHandle( final Handle handle, final SQLMetadataConnector connector, final MetadataStorageTablesConfig dbTables, + final SegmentsMetadataManagerConfig managerConfig, final ObjectMapper jsonMapper ) { - return new SqlSegmentsMetadataQuery(handle, connector, dbTables, jsonMapper); + return new SqlSegmentsMetadataQuery(handle, connector, dbTables, managerConfig, jsonMapper); } /** @@ -902,6 +908,17 @@ public int markSegmentsUnused( public boolean markSegmentAsUsed(SegmentId segmentId, DateTime updateTime) { + final List targetSegments = + retrieveSegmentsById(segmentId.getDataSource(), Set.of(segmentId)); + if (targetSegments.isEmpty()) { + // Invalid segment ID + return false; + } else if (Boolean.TRUE.equals(targetSegments.getFirst().getUsed())) { + // Segment is already marked as used, avoid another DB call + return false; + } + + validateSegmentsForMarkingAsUsed(targetSegments); return markSegments(Set.of(segmentId), true, updateTime) > 0; } @@ -928,7 +945,6 @@ private int markNonOvershadowedSegmentsAsUsedInternal( DateTime updateTime ) { - final List unusedSegments = new ArrayList<>(); final SegmentTimeline timeline = new SegmentTimeline(); final List intervals = @@ -942,12 +958,13 @@ private int markNonOvershadowedSegmentsAsUsedInternal( throw new RuntimeException(e); } - try (final CloseableIterator iterator = - retrieveUnusedSegments(dataSourceName, intervals, versions, null, null, null, null)) { + final List unusedSegments = new ArrayList<>(); + try (final CloseableIterator iterator = + retrieveUnusedSegmentsPlus(dataSourceName, intervals, versions, null, null, null, null)) { while (iterator.hasNext()) { - final DataSegment dataSegment = iterator.next(); - timeline.add(dataSegment); - unusedSegments.add(dataSegment); + final DataSegmentPlus segmentPlus = iterator.next(); + timeline.add(segmentPlus.getDataSegment()); + unusedSegments.add(segmentPlus); } } catch (IOException e) { @@ -958,18 +975,18 @@ private int markNonOvershadowedSegmentsAsUsedInternal( } private int markNonOvershadowedSegmentsAsUsed( - List unusedSegments, + List unusedSegments, SegmentTimeline timeline, DateTime updateTime ) { - Set nonOvershadowedSegments = + final Map nonOvershadowedSegments = unusedSegments.stream() - .filter(segment -> !timeline.isOvershadowed(segment)) - .map(DataSegment::getId) - .collect(Collectors.toSet()); + .filter(segment -> !timeline.isOvershadowed(segment.getDataSegment())) + .collect(Collectors.toMap(segment -> segment.getDataSegment().getId(), Function.identity())); - return markSegmentsAsUsed(nonOvershadowedSegments, updateTime); + validateSegmentsForMarkingAsUsed(nonOvershadowedSegments.values()); + return markSegmentsAsUsed(nonOvershadowedSegments.keySet(), updateTime); } public int markNonOvershadowedSegmentsAsUsed( @@ -978,13 +995,15 @@ public int markNonOvershadowedSegmentsAsUsed( final DateTime updateTime ) { - final List unusedSegments = retrieveUnusedSegments(dataSource, segmentIds); + final List unusedSegments = retrieveUnusedSegmentsPlus(dataSource, segmentIds); final List unusedSegmentsIntervals = JodaUtils.condenseIntervals( - unusedSegments.stream().map(DataSegment::getInterval).collect(Collectors.toList()) + unusedSegments.stream().map(s -> s.getDataSegment().getInterval()).toList() ); // Create a timeline with all used and unused segments in this interval - final SegmentTimeline timeline = SegmentTimeline.forSegments(unusedSegments); + final SegmentTimeline timeline = SegmentTimeline.forSegments( + unusedSegments.stream().map(DataSegmentPlus::getDataSegment).toList() + ); try (CloseableIterator usedSegmentsOverlappingUnusedSegmentsIntervals = retrieveUsedSegments(dataSource, unusedSegmentsIntervals)) { @@ -997,7 +1016,44 @@ public int markNonOvershadowedSegmentsAsUsed( return markNonOvershadowedSegmentsAsUsed(unusedSegments, timeline, updateTime); } - private List retrieveUnusedSegments( + /** + * Checks that all the segments were last updated within the kill buffer period. + * If any segment was updated earlier than that, an exception is thrown so that + * none of the segments are updated to ensure atomicity. + */ + private void validateSegmentsForMarkingAsUsed(Collection segments) + { + if (!managerConfig.getKillUnused().isEnabled()) { + // Do not verify the buffer period if embedded kill tasks are not enabled + return; + } + + final Period bufferPeriod = managerConfig.getKillUnused().getBufferPeriod(); + final DateTime minAllowedUpdateTime = DateTimes.nowUtc().minus( + managerConfig.getKillUnused().getBufferPeriod() + ); + + final List expiredSegmentIds = segments.stream().filter( + s -> s.getUsedStatusLastUpdatedDate() != null + && !s.getUsedStatusLastUpdatedDate().isAfter(minAllowedUpdateTime) + ).map(s -> s.getDataSegment().getId()).toList(); + + if (!expiredSegmentIds.isEmpty()) { + throw DruidException.forPersona(DruidException.Persona.OPERATOR) + .ofCategory(DruidException.Category.CONFLICT) + .build( + "Segment IDs[%s] cannot be marked as used since" + + " they were last updated more than [%s] ago and" + + " are now eligible for permanent deletion." + + " Increase the value of runtime property" + + " ['druid.manager.segments.killUnused.bufferPeriod']" + + " to allow updating these segment IDs.", + expiredSegmentIds, bufferPeriod + ); + } + } + + private List retrieveUnusedSegmentsPlus( final String dataSource, final Set segmentIds ) @@ -1005,12 +1061,11 @@ private List retrieveUnusedSegments( final List retrievedSegments = retrieveSegmentsById(dataSource, segmentIds); final Set unknownSegmentIds = new HashSet<>(segmentIds); - final List unusedSegments = new ArrayList<>(); + final List unusedSegments = new ArrayList<>(); for (DataSegmentPlus entry : retrievedSegments) { - final DataSegment segment = entry.getDataSegment(); - unknownSegmentIds.remove(segment.getId()); + unknownSegmentIds.remove(entry.getDataSegment().getId()); if (Boolean.FALSE.equals(entry.getUsed())) { - unusedSegments.add(segment); + unusedSegments.add(entry); } } @@ -1574,7 +1629,9 @@ private CloseableIterator retrieveSegmentsPlus( @Nullable final DateTime maxUsedStatusLastUpdatedTime ) { - if (intervals.isEmpty() || intervals.size() <= MAX_INTERVALS_PER_BATCH) { + if (versions != null && versions.isEmpty()) { + return CloseableIterators.withEmptyBaggage(Collections.emptyIterator()); + } else if (intervals.isEmpty() || intervals.size() <= MAX_INTERVALS_PER_BATCH) { return CloseableIterators.withEmptyBaggage( retrieveSegmentsPlusInIntervalsBatch(dataSource, intervals, versions, matchMode, used, limit, lastSegmentId, sortOrder, maxUsedStatusLastUpdatedTime) ); diff --git a/server/src/main/java/org/apache/druid/metadata/UnusedSegmentKillerConfig.java b/server/src/main/java/org/apache/druid/metadata/UnusedSegmentKillerConfig.java index b2deed50438e..98fb77ab9329 100644 --- a/server/src/main/java/org/apache/druid/metadata/UnusedSegmentKillerConfig.java +++ b/server/src/main/java/org/apache/druid/metadata/UnusedSegmentKillerConfig.java @@ -23,7 +23,9 @@ import com.fasterxml.jackson.annotation.JsonProperty; import org.apache.druid.common.config.Configs; import org.apache.druid.error.InvalidInput; +import org.apache.druid.java.util.common.DateTimes; import org.apache.druid.java.util.common.logger.Logger; +import org.joda.time.DateTime; import org.joda.time.Period; import javax.annotation.Nullable; @@ -44,6 +46,15 @@ public class UnusedSegmentKillerConfig */ public static final int DEFAULT_MAX_SEGMENTS_TO_KILL = 200_000; + public static final Period DEFAULT_BUFFER_PERIOD = Period.days(30); + + /** + * Grace period as used in {@link #getMaxUpdatedTimeOfKillableSegment()}. + * A grace period of 1 hour is adequate to allow any ongoing segment update + * operations to finish. + */ + public static final Period GRACE_PERIOD = Period.hours(1); + @JsonProperty("enabled") private final boolean enabled; @@ -65,7 +76,7 @@ public UnusedSegmentKillerConfig( ) { this.enabled = Configs.valueOrDefault(enabled, false); - this.bufferPeriod = Configs.valueOrDefault(bufferPeriod, Period.days(30)); + this.bufferPeriod = Configs.valueOrDefault(bufferPeriod, DEFAULT_BUFFER_PERIOD); this.maxSegmentsToKill = Configs.valueOrDefault(maxSegmentsToKill, DEFAULT_MAX_SEGMENTS_TO_KILL); if (this.maxSegmentsToKill > DEFAULT_MAX_SEGMENTS_TO_KILL) { @@ -95,13 +106,29 @@ public UnusedSegmentKillerConfig( /** * Period for which segments are retained even after being marked as unused. - * Default value is 30 days. + * Default value is {@link #DEFAULT_BUFFER_PERIOD}. */ public Period getBufferPeriod() { return bufferPeriod; } + /** + * Maximum value for the updated time of a segment that makes it eligible for + * kill. A segment becomes eligible if it has been unused for at least the + * {@link #getBufferPeriod()}. After this period, the segment cannot be marked + * as used again anymore. Since marking a non-overshadowed segment as used can + * be a slow operation (due to the requirement to build the entire timeline + * and then identify non-overshadowed segments), a {@link #GRACE_PERIOD} is + * added to the buffer period. This helps avoid any unexpected behaviour in + * case a slow update operation is started right at the boundary of the buffer + * period, and a kill task is launched right after. + */ + public DateTime getMaxUpdatedTimeOfKillableSegment() + { + return DateTimes.nowUtc().minus(bufferPeriod.plus(GRACE_PERIOD)); + } + /** * Period dictating the frequency at which the unused segment killer duty * should be run. This config is for testing only and SHOULD NOT be used in diff --git a/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataReadTransaction.java b/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataReadTransaction.java index 924a8c3c54b8..af835537ff18 100644 --- a/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataReadTransaction.java +++ b/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataReadTransaction.java @@ -38,7 +38,9 @@ public interface SegmentMetadataReadTransaction Handle getHandle(); /** - * @return SQL tool to read or update the metadata store directly. + * SQL tool to read or update the metadata store directly, without affecting + * the cache. Use this when performing a no-cache operation inside a transaction + * that otherwise uses the cache for some operations. */ SqlSegmentsMetadataQuery noCacheSql(); diff --git a/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataTransactionFactory.java b/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataTransactionFactory.java index e048e10970c6..2daa9f906a19 100644 --- a/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataTransactionFactory.java +++ b/server/src/main/java/org/apache/druid/metadata/segment/SegmentMetadataTransactionFactory.java @@ -19,6 +19,10 @@ package org.apache.druid.metadata.segment; +import org.apache.druid.metadata.SqlSegmentsMetadataQuery; + +import java.util.function.Function; + /** * Factory for {@link SegmentMetadataTransaction}s. */ @@ -41,4 +45,21 @@ T inReadWriteDatasourceTransaction( String dataSource, SegmentMetadataTransaction.Callback callback ); + + + /** + * Performs a read-only transaction which queries the metadata store. + */ + T inReadOnlyNoCacheTransaction( + Function sqlQuery + ); + + /** + * Performs a write transaction which updates the metadata store directly, + * and does not affect the cache. + */ + T inReadWriteNoCacheTransaction( + Function sqlUpdate + ); + } diff --git a/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataReadOnlyTransactionFactory.java b/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataReadOnlyTransactionFactory.java index 0d0d1e42a1e2..40c8e6ef6a23 100644 --- a/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataReadOnlyTransactionFactory.java +++ b/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataReadOnlyTransactionFactory.java @@ -24,9 +24,13 @@ import org.apache.druid.error.DruidException; import org.apache.druid.metadata.MetadataStorageTablesConfig; import org.apache.druid.metadata.SQLMetadataConnector; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; +import org.apache.druid.metadata.SqlSegmentsMetadataQuery; import org.skife.jdbi.v2.Handle; import org.skife.jdbi.v2.TransactionStatus; +import java.util.function.Function; + /** * Factory for read-only {@link SegmentMetadataTransaction}s that always read * directly from the metadata store and never from the {@code SegmentMetadataCache}. @@ -40,17 +44,20 @@ public class SqlSegmentMetadataReadOnlyTransactionFactory implements SegmentMeta private final ObjectMapper jsonMapper; private final MetadataStorageTablesConfig tablesConfig; + private final SegmentsMetadataManagerConfig managerConfig; private final SQLMetadataConnector connector; @Inject public SqlSegmentMetadataReadOnlyTransactionFactory( ObjectMapper jsonMapper, MetadataStorageTablesConfig tablesConfig, + SegmentsMetadataManagerConfig managerConfig, SQLMetadataConnector connector ) { this.jsonMapper = jsonMapper; this.tablesConfig = tablesConfig; + this.managerConfig = managerConfig; this.connector = connector; } @@ -90,6 +97,22 @@ public T inReadWriteDatasourceTransaction( throw DruidException.defensive("Only Overlord can perform write transactions on segment metadata."); } + @Override + public T inReadOnlyNoCacheTransaction(Function sqlQuery) + { + return connector.retryReadOnlyTransaction( + (handle, status) -> sqlQuery.apply(createSqlQueryForTransaction(handle)), + getQuietRetries(), + getMaxRetries() + ); + } + + @Override + public T inReadWriteNoCacheTransaction(Function sqlUpdate) + { + throw DruidException.defensive("Only Overlord can perform write transactions on segment metadata."); + } + protected SegmentMetadataTransaction createSqlTransaction( String dataSource, Handle handle, @@ -98,7 +121,9 @@ protected SegmentMetadataTransaction createSqlTransaction( { return new SqlSegmentMetadataTransaction( dataSource, - handle, transactionStatus, connector, tablesConfig, jsonMapper + handle, + createSqlQueryForTransaction(handle), + transactionStatus, connector, tablesConfig, jsonMapper ); } @@ -111,4 +136,9 @@ protected T executeReadAndClose( return callback.inTransaction(transaction); } } + + protected SqlSegmentsMetadataQuery createSqlQueryForTransaction(Handle handle) + { + return SqlSegmentsMetadataQuery.forHandle(handle, connector, tablesConfig, managerConfig, jsonMapper); + } } diff --git a/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransaction.java b/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransaction.java index 4e144433260f..e1ae6f2a9577 100644 --- a/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransaction.java +++ b/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransaction.java @@ -75,6 +75,7 @@ class SqlSegmentMetadataTransaction implements SegmentMetadataTransaction SqlSegmentMetadataTransaction( String dataSource, Handle handle, + SqlSegmentsMetadataQuery query, TransactionStatus transactionStatus, SQLMetadataConnector connector, MetadataStorageTablesConfig dbTables, @@ -87,7 +88,7 @@ class SqlSegmentMetadataTransaction implements SegmentMetadataTransaction this.dbTables = dbTables; this.jsonMapper = jsonMapper; this.transactionStatus = transactionStatus; - this.query = SqlSegmentsMetadataQuery.forHandle(handle, connector, dbTables, jsonMapper); + this.query = query; } @Override diff --git a/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransactionFactory.java b/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransactionFactory.java index 26083e44b678..38d81aa7ef7e 100644 --- a/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransactionFactory.java +++ b/server/src/main/java/org/apache/druid/metadata/segment/SqlSegmentMetadataTransactionFactory.java @@ -28,10 +28,14 @@ import org.apache.druid.java.util.emitter.service.ServiceMetricEvent; import org.apache.druid.metadata.MetadataStorageTablesConfig; import org.apache.druid.metadata.SQLMetadataConnector; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; +import org.apache.druid.metadata.SqlSegmentsMetadataQuery; import org.apache.druid.metadata.segment.cache.Metric; import org.apache.druid.metadata.segment.cache.SegmentMetadataCache; import org.apache.druid.query.DruidMetrics; +import java.util.function.Function; + /** * Factory for {@link SegmentMetadataTransaction}s. If the * {@link SegmentMetadataCache} is enabled and ready, the transaction may @@ -65,10 +69,11 @@ public SqlSegmentMetadataTransactionFactory( SQLMetadataConnector connector, @IndexingService DruidLeaderSelector leaderSelector, SegmentMetadataCache segmentMetadataCache, + SegmentsMetadataManagerConfig managerConfig, ServiceEmitter emitter ) { - super(jsonMapper, tablesConfig, connector); + super(jsonMapper, tablesConfig, managerConfig, connector); this.connector = connector; this.leaderSelector = leaderSelector; this.segmentMetadataCache = segmentMetadataCache; @@ -142,6 +147,16 @@ public T inReadWriteDatasourceTransaction( ); } + @Override + public T inReadWriteNoCacheTransaction(Function sqlUpdate) + { + return connector.retryTransaction( + (handle, status) -> sqlUpdate.apply(createSqlQueryForTransaction(handle)), + getQuietRetries(), + getMaxRetries() + ); + } + private T executeWriteAndClose( SegmentMetadataTransaction transaction, SegmentMetadataTransaction.Callback callback diff --git a/server/src/main/java/org/apache/druid/metadata/segment/cache/HeapMemorySegmentMetadataCache.java b/server/src/main/java/org/apache/druid/metadata/segment/cache/HeapMemorySegmentMetadataCache.java index d9d2cfce1596..6b62e6085215 100644 --- a/server/src/main/java/org/apache/druid/metadata/segment/cache/HeapMemorySegmentMetadataCache.java +++ b/server/src/main/java/org/apache/druid/metadata/segment/cache/HeapMemorySegmentMetadataCache.java @@ -134,6 +134,7 @@ private enum CacheState private final Duration pollDuration; private final UsageMode cacheMode; private final MetadataStorageTablesConfig tablesConfig; + private final SegmentsMetadataManagerConfig managerConfig; private final SQLMetadataConnector connector; private final boolean useSchemaCache; @@ -180,8 +181,9 @@ public HeapMemorySegmentMetadataCache( ) { this.jsonMapper = jsonMapper; - this.cacheMode = config.get().getCacheUsageMode(); - this.pollDuration = config.get().getPollDuration().toStandardDuration(); + this.managerConfig = config.get(); + this.cacheMode = managerConfig.getCacheUsageMode(); + this.pollDuration = managerConfig.getPollDuration().toStandardDuration(); this.tablesConfig = tablesConfig.get(); this.useSchemaCache = segmentSchemaCache.isEnabled(); this.segmentSchemaCache = segmentSchemaCache; @@ -750,7 +752,7 @@ private T query(Function sqlFunction) return inReadOnlyTransaction( (handle, status) -> sqlFunction.apply( SqlSegmentsMetadataQuery - .forHandle(handle, connector, tablesConfig, jsonMapper) + .forHandle(handle, connector, tablesConfig, managerConfig, jsonMapper) ) ); } @@ -780,7 +782,7 @@ private void retrieveRequiredUsedSegments( try ( CloseableIterator iterator = SqlSegmentsMetadataQuery - .forHandle(handle, connector, tablesConfig, jsonMapper) + .forHandle(handle, connector, tablesConfig, managerConfig, jsonMapper) .retrieveSegmentsByIdIterator(dataSource, segmentIdsToRefresh, useSchemaCache) ) { iterator.forEachRemaining(summary.usedSegments::add); diff --git a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorMarkUsedTest.java b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorMarkUsedTest.java index dadad5e4f1ba..c807596d8828 100644 --- a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorMarkUsedTest.java +++ b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorMarkUsedTest.java @@ -80,6 +80,7 @@ public void setup() derbyConnector, new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ) { diff --git a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorReadOnlyTest.java b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorReadOnlyTest.java index d68da314ad16..0ad44376819d 100644 --- a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorReadOnlyTest.java +++ b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorReadOnlyTest.java @@ -162,6 +162,7 @@ private IndexerSQLMetadataStorageCoordinator createStorageCoordinator( transactionFactory = new SqlSegmentMetadataReadOnlyTransactionFactory( mapper, derbyConnectorRule.metadataTablesConfigSupplier().get(), + new SegmentsMetadataManagerConfig(null, null, null), derbyConnector ); } else { @@ -171,6 +172,7 @@ private IndexerSQLMetadataStorageCoordinator createStorageCoordinator( derbyConnector, leaderSelector, segmentMetadataCache, + new SegmentsMetadataManagerConfig(null, cacheMode, null), emitter ); } 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 45f02c3713aa..a6a4f8132c66 100644 --- a/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorTest.java +++ b/server/src/test/java/org/apache/druid/metadata/IndexerSQLMetadataStorageCoordinatorTest.java @@ -25,6 +25,7 @@ import com.google.common.collect.Iterables; import org.apache.druid.common.utils.IdUtils; import org.apache.druid.data.input.StringTuple; +import org.apache.druid.error.DruidException; import org.apache.druid.error.DruidExceptionMatcher; import org.apache.druid.error.ExceptionMatcher; import org.apache.druid.indexer.partitions.DynamicPartitionsSpec; @@ -92,7 +93,6 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.Parameterized; -import org.skife.jdbi.v2.exceptions.CallbackFailedException; import java.io.IOException; import java.nio.charset.StandardCharsets; @@ -206,6 +206,7 @@ public void setUp() derbyConnector, leaderSelector, segmentMetadataCache, + new SegmentsMetadataManagerConfig(null, cacheMode, null), emitter ) { @@ -599,9 +600,9 @@ public void testCommitReplaceSegments_partiallyOverlappingPendingSegmentUnsuppor replacingSegments.add(segment); } - Assert.assertFalse( - coordinator.commitReplaceSegments(replacingSegments, ImmutableSet.of(replaceLock), null) - .isSuccess() + Assert.assertThrows( + DruidException.class, + () -> coordinator.commitReplaceSegments(replacingSegments, ImmutableSet.of(replaceLock), null) ); } @@ -1590,7 +1591,7 @@ public void testRetrieveUnusedSegmentsUsingSingleIntervalAndLimitInRange() ); Assert.assertEquals(requestedLimit, actualUnusedSegments.size()); - Assert.assertTrue(actualUnusedSegments.containsAll(segments.stream().limit(requestedLimit).collect(Collectors.toList()))); + Assert.assertTrue(actualUnusedSegments.containsAll(segments.stream().limit(requestedLimit).toList())); } @Test @@ -3597,16 +3598,14 @@ public void test_allocateCommitDelete_createsFreshVersion_uptoMaxAllowedRetries( // Verify that the next attempt fails MatcherAssert.assertThat( Assert.assertThrows( - CallbackFailedException.class, + DruidException.class, () -> allocatePendingSegmentForAppendTask(wiki, firstOfJan23, IdUtils.getRandomId()) ), - ExceptionMatcher.of(CallbackFailedException.class).expectRootCause( - DruidExceptionMatcher.internalServerError().expectMessageIs( - "Could not allocate segment" - + "[wiki_2023-01-01T00:00:00.000Z_2023-01-02T00:00:00.000Z_1970-01-01T00:00:00.000Z]" - + " as there are too many clashing unused versions(upto [1970-01-01T00:00:00.000ZSSSSSSSSSS])" - + " in the interval. Kill the old unused versions to proceed." - ) + DruidExceptionMatcher.internalServerError().expectMessageIs( + "Could not allocate segment" + + "[wiki_2023-01-01T00:00:00.000Z_2023-01-02T00:00:00.000Z_1970-01-01T00:00:00.000Z]" + + " as there are too many clashing unused versions(upto [1970-01-01T00:00:00.000ZSSSSSSSSSS])" + + " in the interval. Kill the old unused versions to proceed." ) ); } diff --git a/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorSchemaPersistenceTest.java b/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorSchemaPersistenceTest.java index ba64c709fb18..ec37be03ac47 100644 --- a/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorSchemaPersistenceTest.java +++ b/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorSchemaPersistenceTest.java @@ -101,6 +101,7 @@ public void setUp() derbyConnector, new TestDruidLeaderSelector(), NoopSegmentMetadataCache.instance(), + new SegmentsMetadataManagerConfig(null, null, null), NoopServiceEmitter.instance() ); coordinator = new IndexerSQLMetadataStorageCoordinator( diff --git a/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorTestBase.java b/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorTestBase.java index be0120162cf0..845be8c48243 100644 --- a/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorTestBase.java +++ b/server/src/test/java/org/apache/druid/metadata/IndexerSqlMetadataStorageCoordinatorTestBase.java @@ -359,6 +359,7 @@ protected ImmutableList retrieveUnusedSegments( handle, derbyConnector, tablesConfig, + new SegmentsMetadataManagerConfig(null, null, null), mapper ) .retrieveUnusedSegments( @@ -385,10 +386,11 @@ protected ImmutableList retrieveUnusedSegmentsPlus( MetadataStorageTablesConfig tablesConfig ) { + final SegmentsMetadataManagerConfig managerConfig = new SegmentsMetadataManagerConfig(null, null, null); return derbyConnector.inReadOnlyTransaction( (handle, status) -> { try (final CloseableIterator iterator = - SqlSegmentsMetadataQuery.forHandle(handle, derbyConnector, tablesConfig, mapper) + SqlSegmentsMetadataQuery.forHandle(handle, derbyConnector, tablesConfig, managerConfig, mapper) .retrieveUnusedSegmentsPlus( TestDataSource.WIKI, intervals, diff --git a/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataManagerTestBase.java b/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataManagerTestBase.java index bedc41f11e81..f7f52077987e 100644 --- a/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataManagerTestBase.java +++ b/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataManagerTestBase.java @@ -124,8 +124,9 @@ public static int markSegmentsAsUnused( DateTime updateTime ) { + final SegmentsMetadataManagerConfig managerConfig = new SegmentsMetadataManagerConfig(null, null, null); return connector.retryWithHandle( - handle -> SqlSegmentsMetadataQuery.forHandle(handle, connector, storageConfig, jsonMapper) + handle -> SqlSegmentsMetadataQuery.forHandle(handle, connector, storageConfig, managerConfig, jsonMapper) .markSegmentsAsUnused(segmentIds, updateTime) ); } diff --git a/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataQueryTest.java b/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataQueryTest.java index bd480febbdfb..b0e8c61cbec4 100644 --- a/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataQueryTest.java +++ b/server/src/test/java/org/apache/druid/metadata/SqlSegmentsMetadataQueryTest.java @@ -22,9 +22,12 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.collect.ImmutableSet; import org.apache.druid.data.input.impl.DimensionsSpec; +import org.apache.druid.error.DruidException; +import org.apache.druid.error.DruidExceptionMatcher; import org.apache.druid.indexer.partitions.DynamicPartitionsSpec; import org.apache.druid.java.util.common.DateTimes; import org.apache.druid.java.util.common.Intervals; +import org.apache.druid.java.util.common.StringUtils; import org.apache.druid.java.util.common.granularity.Granularities; import org.apache.druid.java.util.common.parsers.CloseableIterator; import org.apache.druid.metadata.segment.cache.IndexingStateRecord; @@ -37,6 +40,7 @@ import org.apache.druid.timeline.CompactionState; import org.apache.druid.timeline.DataSegment; import org.apache.druid.timeline.SegmentId; +import org.hamcrest.MatcherAssert; import org.joda.time.DateTime; import org.joda.time.Interval; import org.joda.time.Period; @@ -50,6 +54,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.function.BiFunction; import java.util.function.Function; import java.util.stream.Collectors; @@ -308,7 +313,13 @@ private T read(Function function) final MetadataStorageTablesConfig tablesConfig = derbyConnectorRule.metadataTablesConfigSupplier().get(); return connector.inReadOnlyTransaction( (handle, status) -> function.apply( - SqlSegmentsMetadataQuery.forHandle(handle, connector, tablesConfig, TestHelper.JSON_MAPPER) + SqlSegmentsMetadataQuery.forHandle( + handle, + connector, + tablesConfig, + new SegmentsMetadataManagerConfig(null, null, null), + TestHelper.JSON_MAPPER + ) ) ); } @@ -324,7 +335,13 @@ private Set readAsSet(Function { final SqlSegmentsMetadataQuery query = - SqlSegmentsMetadataQuery.forHandle(handle, connector, tablesConfig, TestHelper.JSON_MAPPER); + SqlSegmentsMetadataQuery.forHandle( + handle, + connector, + tablesConfig, + new SegmentsMetadataManagerConfig(null, null, null), + TestHelper.JSON_MAPPER + ); try (CloseableIterator iterator = iterableReader.apply(query)) { return ImmutableSet.copyOf(iterator); @@ -336,12 +353,33 @@ private Set readAsSet(Function T update(Function function) + { + return updateWithConfig( + function, + new SegmentsMetadataManagerConfig(null, null, null) + ); + } + + /** + * Executes an update using a {@link SqlSegmentsMetadataQuery} object initialized + * with the given {@link SegmentsMetadataManagerConfig}. + */ + private T updateWithConfig( + Function function, + SegmentsMetadataManagerConfig managerConfig + ) { final DerbyConnector connector = derbyConnectorRule.getConnector(); final MetadataStorageTablesConfig tablesConfig = derbyConnectorRule.metadataTablesConfigSupplier().get(); return connector.retryWithHandle( handle -> function.apply( - SqlSegmentsMetadataQuery.forHandle(handle, connector, tablesConfig, TestHelper.JSON_MAPPER) + SqlSegmentsMetadataQuery.forHandle( + handle, + connector, + tablesConfig, + managerConfig, + TestHelper.JSON_MAPPER + ) ) ); } @@ -375,6 +413,133 @@ private static Set getIds(Set segments) return segments.stream().map(DataSegment::getId).collect(Collectors.toSet()); } + // ==================== Kill Buffer Period Tests ==================== + + @Test + public void test_markSegmentAsUsed_throwsIfExpiredAndKillEnabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).minusMinutes(1); + verifyMarkAsUsedThrowsConflictException( + (sql, segment) -> sql.markSegmentAsUsed(segment.getId(), DateTimes.nowUtc()), + markedUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(true, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markNonOvershadowedSegmentsAsUsed_throwsIfExpiredAndKillEnabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).minusMinutes(1); + verifyMarkAsUsedThrowsConflictException( + (sql, segment) -> sql.markNonOvershadowedSegmentsAsUsed( + TestDataSource.WIKI, + Set.of(segment.getId()), + DateTimes.nowUtc() + ), + markedUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(true, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markAllNonOvershadowedSegmentsAsUsed_throwsIfExpiredAndKillEnabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).minusMinutes(1); + verifyMarkAsUsedThrowsConflictException( + (sql, segment) -> sql.markAllNonOvershadowedSegmentsAsUsed( + segment.getDataSource(), + DateTimes.nowUtc() + ), + markedUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(true, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markSegmentAsUsed_succeedsIfRecentlyUpdatedAndKillEnabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedAsUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).plusMinutes(1); + verifyMarkAsUsedSucceeds( + (sql, segment) -> sql.markSegmentAsUsed(segment.getId(), DateTimes.nowUtc()) ? 1 : 0, + markedAsUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(true, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markSegmentsAsUsed_succeedsIfRecentlyUpdatedAndKillEnabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedAsUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).plusMinutes(1); + verifyMarkAsUsedSucceeds( + (sql, segment) -> sql.markSegmentsAsUsed(Set.of(segment.getId()), DateTimes.nowUtc()), + markedAsUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(true, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markAllNonOvershadowedSegmentsAsUsed_succeedsIfRecentlyUpdatedAndKillEnabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedAsUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).plusMinutes(1); + verifyMarkAsUsedSucceeds( + (sql, segment) -> sql.markAllNonOvershadowedSegmentsAsUsed( + segment.getDataSource(), + DateTimes.nowUtc() + ), + markedAsUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(true, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markNonOvershadowedSegmentsAsUsed_succeedsIfExpiredButKillDisabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedAsUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).minusDays(1); + verifyMarkAsUsedSucceeds( + (sql, segment) -> sql.markSegmentsAsUsed(Set.of(segment.getId()), DateTimes.nowUtc()), + markedAsUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(false, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markSegmentsAsUsed_succeedsIfExpiredButKillDisabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedAsUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).minusDays(1); + verifyMarkAsUsedSucceeds( + (sql, segment) -> sql.markNonOvershadowedSegmentsAsUsed( + segment.getDataSource(), + Set.of(segment.getId()), + DateTimes.nowUtc() + ), + markedAsUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(false, bufferPeriod, null, null)) + ); + } + + @Test + public void test_markAllNonOvershadowedSegmentsAsUsed_succeedsIfExpiredButKillDisabled() + { + final Period bufferPeriod = Period.days(60); + final DateTime markedAsUnusedAtTime = DateTimes.nowUtc().minus(bufferPeriod).minusDays(1); + verifyMarkAsUsedSucceeds( + (sql, segment) -> sql.markAllNonOvershadowedSegmentsAsUsed( + segment.getDataSource(), + DateTimes.nowUtc() + ), + markedAsUnusedAtTime, + createManagerConfig(new UnusedSegmentKillerConfig(false, bufferPeriod, null, null)) + ); + } + // ==================== Indexing State Tests ==================== @Test @@ -717,4 +882,76 @@ private void markIndexingStateAsUnused(String fingerprint) return null; }); } + + /** + * Marks a single segment as unused, sets the used_status_last_updated equal + * to the {@param markedAsUnusedAtTime} and then tries to mark it as used with + * the given function. + */ + private void verifyMarkAsUsedSucceeds( + BiFunction markAsUsedFunction, + DateTime markedAsUnusedAtTime, + SegmentsMetadataManagerConfig managerConfig + ) + { + final DataSegment segment = WIKI_SEGMENTS_2X5D.getFirst(); + + // Mark segments as unused with the given used_status_last_updated time + update(sql -> sql.markSegmentsAsUnused(Set.of(segment.getId()), markedAsUnusedAtTime)); + updateUsedStatusLastUpdated(segment.getId(), markedAsUnusedAtTime); + + final int numUpdatedRows = updateWithConfig(sql -> markAsUsedFunction.apply(sql, segment), managerConfig); + Assert.assertEquals(1, numUpdatedRows); + Assert.assertTrue(retrieveAllUsedSegments().contains(segment)); + } + + /** + * Marks a single segment as unused, sets the used_status_last_updated equal + * to the {@param markedAsUnusedAtTime}, and then tries to mark it as used + * with the given function. + */ + private void verifyMarkAsUsedThrowsConflictException( + BiFunction markAsUsedFunction, + DateTime markedAsUnusedAtTime, + SegmentsMetadataManagerConfig managerConfig + ) + { + final DataSegment segment = WIKI_SEGMENTS_2X5D.getFirst(); + + // Mark segment as unused with an old update time (outside buffer period) + updateUsedStatusLastUpdated(segment.getId(), markedAsUnusedAtTime); + update(sql -> sql.markSegmentsAsUnused(Set.of(segment.getId()), markedAsUnusedAtTime)); + + // Verify that the mark as used operation fails with a CONFLICT DruidException + MatcherAssert.assertThat( + Assert.assertThrows( + DruidException.class, + () -> updateWithConfig(sql -> markAsUsedFunction.apply(sql, segment), managerConfig) + ), + DruidExceptionMatcher.conflict().expectMessageIs( + StringUtils.format( + "Segment IDs[[%s]]" + + " cannot be marked as used since they were last updated more than [%s]" + + " ago and are now eligible for permanent deletion. Increase the value" + + " of runtime property ['druid.manager.segments.killUnused.bufferPeriod']" + + " to allow updating these segment IDs.", + segment.getId(), + managerConfig.getKillUnused().getBufferPeriod() + ) + ) + ); + } + + /** + * Updates the used_status_last_updated column for the given segment. + */ + private void updateUsedStatusLastUpdated(SegmentId segmentId, DateTime updateTime) + { + derbyConnectorRule.segments().updateUsedStatusLastUpdated(segmentId.toString(), updateTime); + } + + private static SegmentsMetadataManagerConfig createManagerConfig(UnusedSegmentKillerConfig killerConfig) + { + return new SegmentsMetadataManagerConfig(null, null, killerConfig); + } } diff --git a/server/src/test/java/org/apache/druid/server/coordinator/duty/KillUnusedSegmentsTest.java b/server/src/test/java/org/apache/druid/server/coordinator/duty/KillUnusedSegmentsTest.java index 37501cb849cc..8fbd48ec28f9 100644 --- a/server/src/test/java/org/apache/druid/server/coordinator/duty/KillUnusedSegmentsTest.java +++ b/server/src/test/java/org/apache/druid/server/coordinator/duty/KillUnusedSegmentsTest.java @@ -39,8 +39,11 @@ import org.apache.druid.metadata.IndexerSQLMetadataStorageCoordinator; import org.apache.druid.metadata.MetadataStorageTablesConfig; import org.apache.druid.metadata.SQLMetadataConnector; +import org.apache.druid.metadata.SegmentsMetadataManagerConfig; import org.apache.druid.metadata.SqlSegmentsMetadataManagerTestBase; import org.apache.druid.metadata.TestDerbyConnector; +import org.apache.druid.metadata.segment.SegmentMetadataTransactionFactory; +import org.apache.druid.metadata.segment.SqlSegmentMetadataReadOnlyTransactionFactory; import org.apache.druid.rpc.indexing.NoopOverlordClient; import org.apache.druid.segment.TestHelper; import org.apache.druid.segment.metadata.CentralizedDatasourceSchemaConfig; @@ -112,8 +115,14 @@ public class KillUnusedSegmentsTest public void setup() { connector = derbyConnectorRule.getConnector(); + final SegmentMetadataTransactionFactory transactionFactory = new SqlSegmentMetadataReadOnlyTransactionFactory( + TestHelper.JSON_MAPPER, + derbyConnectorRule.metadataTablesConfigSupplier().get(), + new SegmentsMetadataManagerConfig(null, null, null), + connector + ); storageCoordinator = new IndexerSQLMetadataStorageCoordinator( - null, + transactionFactory, TestHelper.JSON_MAPPER, derbyConnectorRule.metadataTablesConfigSupplier().get(), connector, diff --git a/services/src/main/java/org/apache/druid/guice/MetadataManagerModule.java b/services/src/main/java/org/apache/druid/guice/MetadataManagerModule.java index 46fc9c835ac6..e9dff3b35c9f 100644 --- a/services/src/main/java/org/apache/druid/guice/MetadataManagerModule.java +++ b/services/src/main/java/org/apache/druid/guice/MetadataManagerModule.java @@ -138,6 +138,7 @@ public void configure(Binder binder) .in(LazySingleton.class); binder.bind(IndexingStateCache.class).in(LazySingleton.class); } else { + // Non-Overlord nodes (i.e. Coordinator) can only read from metadata store binder.bind(SegmentMetadataTransactionFactory.class) .to(SqlSegmentMetadataReadOnlyTransactionFactory.class) .in(LazySingleton.class);