This is an automated email from the ASF dual-hosted git repository.
kfaraz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new 978de632416 feat: Add new TaskLockType KILL for embedded kill tasks
(#19921)
978de632416 is described below
commit 978de63241646ead94f7aab1ca7dabac89ef9e9b
Author: Kashif Faraz <[email protected]>
AuthorDate: Wed Aug 12 21:51:45 2026 +0530
feat: Add new TaskLockType KILL for embedded kill tasks (#19921)
Changes
---------
- Add new `TaskLockType.KILL`
- This lock can coexist with any other lock type but not with another
`KILL` lock
to ensure that only a single `kill` task is working on a single
datasource-interval
- This lock type can currently be used only by embedded kill tasks so that
the
`UnusedSegmentsKiller` does not skip intervals with active ingestion.
- Inject the `SegmentsMetadataManagerConfig` into
`SqlSegmentsMetadataQuery`.
- Do not allow a segment to be marked as used if it was last updated
earlier than the
buffer period. This ensures that an unused segment that is now eligible for
kill is not
accidentally marked as used while a `kill` task for that segment is in
progress.
- Add methods in `SegmentMetadataTransactionFactory` to provide read-write
transactions which do not use the segment metadata cache. This moves away
the
transaction creation logic from `IndexerSQLMetadataStorageCoordinator` to
the
transaction factory and reduces call sites that create a
`SqlSegmentsMetadataQuery`.
- Add tests
---
.../SqlSegmentsMetadataQueryBenchmark.java | 9 +-
docs/configuration/index.md | 4 +-
docs/data-management/delete.md | 2 +-
.../embedded/indexing/KafkaClusterMetricsTest.java | 18 +-
.../apache/druid/indexing/common/TaskLockType.java | 10 +-
.../common/task/KillUnusedSegmentsTask.java | 11 +-
.../druid/indexing/common/task/TaskMetrics.java | 2 +-
.../druid/indexing/overlord/GlobalTaskLockbox.java | 5 +
.../druid/indexing/overlord/TaskLockbox.java | 92 +++++--
.../overlord/duty/UnusedSegmentsKiller.java | 24 +-
.../indexing/common/actions/TaskActionTestKit.java | 5 +-
.../indexing/common/task/IngestionTestBase.java | 5 +-
.../common/task/concurrent/CommandQueueTask.java | 9 +-
.../concurrent/ConcurrentReplaceAndAppendTest.java | 38 ++-
.../ConcurrentReplaceAndStreamingAppendTest.java | 5 +-
.../indexing/overlord/GlobalTaskLockboxTest.java | 269 ++++++++++++++++++++-
.../overlord/TaskLockBoxConcurrencyTest.java | 2 +
.../indexing/overlord/TaskQueueScaleTest.java | 2 +
.../overlord/duty/UnusedSegmentsKillerTest.java | 25 +-
.../SeekableStreamIndexTaskTestBase.java | 2 +
.../IndexerSQLMetadataStorageCoordinator.java | 46 +---
.../druid/metadata/SQLMetadataConnector.java | 16 +-
.../metadata/SegmentsMetadataManagerConfig.java | 8 +-
.../druid/metadata/SqlSegmentsMetadataQuery.java | 105 ++++++--
.../druid/metadata/UnusedSegmentKillerConfig.java | 31 ++-
.../segment/SegmentMetadataReadTransaction.java | 4 +-
.../segment/SegmentMetadataTransactionFactory.java | 21 ++
...lSegmentMetadataReadOnlyTransactionFactory.java | 32 ++-
.../segment/SqlSegmentMetadataTransaction.java | 3 +-
.../SqlSegmentMetadataTransactionFactory.java | 17 +-
.../cache/HeapMemorySegmentMetadataCache.java | 10 +-
...rSQLMetadataStorageCoordinatorMarkUsedTest.java | 1 +
...rSQLMetadataStorageCoordinatorReadOnlyTest.java | 2 +
.../IndexerSQLMetadataStorageCoordinatorTest.java | 25 +-
...ataStorageCoordinatorSchemaPersistenceTest.java | 1 +
...dexerSqlMetadataStorageCoordinatorTestBase.java | 4 +-
.../SqlSegmentsMetadataManagerTestBase.java | 3 +-
.../metadata/SqlSegmentsMetadataQueryTest.java | 243 ++++++++++++++++++-
.../coordinator/duty/KillUnusedSegmentsTest.java | 11 +-
.../apache/druid/guice/MetadataManagerModule.java | 1 +
40 files changed, 952 insertions(+), 171 deletions(-)
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 aa49884cd0b..13fcd970b36 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.granularity.Granularities;
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 @@ public class SqlSegmentsMetadataQueryBenchmark
return derbyConnector.inReadOnlyTransaction((handle, status) -> {
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<T> iterator = iterableReader.apply(query)) {
return ImmutableSet.copyOf(iterator);
diff --git a/docs/configuration/index.md b/docs/configuration/index.md
index d4898222664..238f0b357fa 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.killUnus [...]
#### Overlord dynamic configuration
diff --git a/docs/data-management/delete.md b/docs/data-management/delete.md
index 799e8b4b8b9..54658097c00 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 bc7385e6bf6..49be8522ca0 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.simulate.KafkaResource;
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 class KafkaClusterMetricsTest extends
EmbeddedClusterTestBase
}
};
+ 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 class KafkaClusterMetricsTest extends
EmbeddedClusterTestBase
.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 class KafkaClusterMetricsTest extends
EmbeddedClusterTestBase
@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 class KafkaClusterMetricsTest extends
EmbeddedClusterTestBase
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 43f79cf2651..964e4ef6268 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 @@ package org.apache.druid.indexing.common;
* 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 e31ac044d89..5c74770c230 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 class KillUnusedSegmentsTask extends
AbstractFixedIntervalTask
// 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 86456d12fdc..71250b2a1af 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 @@ public class 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 8b00739167a..6544eec28e1 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 class 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.
+ * <p>
+ * 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 b36d3e96232..2113e5028dc 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.actions.SegmentAllocateResult;
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 @@ public class TaskLockbox
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<TaskLockPosse> foundPosses = findLockPossesOverlapsInterval(
request.getInterval()
@@ -587,7 +596,7 @@ public class TaskLockbox
final List<TaskLockPosse> 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 class TaskLockbox
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 class TaskLockbox
}
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 class TaskLockbox
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 @@ public class TaskLockbox
final List<TaskLockPosse> 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 @@ public class TaskLockbox
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 @@ public class TaskLockbox
/**
* 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 @@ public class TaskLockbox
|| 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 @@ public class TaskLockbox
/**
* 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 @@ public class TaskLockbox
/**
* 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 @@ public class TaskLockbox
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<TaskLockPosse> 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 @@ public class TaskLockbox
final List<TaskLockPosse> 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 @@ public class TaskLockbox
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 @@ public class TaskLockbox
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 a5f8ffa9dc9..2cfa95004ad 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.client.indexing.IndexingService;
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.TaskMetrics;
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 @@ public class UnusedSegmentsKiller implements OverlordDuty
// Identify intervals with unused segments which are eligible for kill
final Map<DatasourceInterval, Integer> eligibleIntervals =
storageCoordinator.retrieveSomeUnusedSegmentIntervals(
- DateTimes.nowUtc().minus(killConfig.getBufferPeriod()),
+ killConfig.getMaxUpdatedTimeOfKillableSegment(),
MAX_INTERVALS_TO_KILL,
killConfig.getMaxSegmentsToKill()
);
@@ -372,7 +377,7 @@ public class UnusedSegmentsKiller implements OverlordDuty
final EmbeddedKillTask killTask = new EmbeddedKillTask(
taskId,
candidate,
- DateTimes.nowUtc().minus(killConfig.getBufferPeriod())
+ killConfig.getMaxUpdatedTimeOfKillableSegment()
);
final TaskActionClient taskActionClient =
taskActionClientFactory.create(killTask);
@@ -474,13 +479,22 @@ public class UnusedSegmentsKiller implements OverlordDuty
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<DataSegmentPlus> 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 8721148083c..3fc1d9d188a 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 @@ public class TaskActionTestKit extends ExternalResource
implements BeforeEachCal
? 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 @@ public class TaskActionTestKit extends ExternalResource
implements BeforeEachCal
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 090ed20269a..0c66a619faa 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 @@ public abstract class IngestionTestBase extends
InitializedNullHandlingTest
= 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 @@ public abstract class IngestionTestBase extends
InitializedNullHandlingTest
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 5945d1488c3..a9e5607c89d 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 @@ public class CommandQueueTask extends AbstractTask
implements PendingSegmentAllo
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 2aa0895c961..6a11cc30ce4 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 class ConcurrentReplaceAndAppendTest extends
IngestionTestBase
// 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 class ConcurrentReplaceAndAppendTest extends
IngestionTestBase
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 class ConcurrentReplaceAndAppendTest extends
IngestionTestBase
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 class ConcurrentReplaceAndAppendTest extends
IngestionTestBase
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 class ConcurrentReplaceAndAppendTest extends
IngestionTestBase
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 class ConcurrentReplaceAndAppendTest extends
IngestionTestBase
}
// 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 7bc0b668561..57c25a65060 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.base.Throwables;
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 @@ import java.util.concurrent.atomic.AtomicInteger;
* <li>LOCK: Acquisition of a lock on an interval by a replace task</li>
* <li>ALLOCATE: Allocation of a pending segment by an append task</li>
* <li>REPLACE: Commit of segments created by a replace task</li>
- * <li>APPEND: Commit of segments created by an append task</li>
+ * <li>APPEND: Commit of segments created by an append task</li>
* </ul>
*/
public class ConcurrentReplaceAndStreamingAppendTest extends IngestionTestBase
@@ -606,7 +607,7 @@ public class ConcurrentReplaceAndStreamingAppendTest
extends IngestionTestBase
// 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 42e6e365b30..c2f9f25e689 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.ImmutableList;
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.AbstractTask;
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.DerbyMetadataStorageActionHandlerFactory;
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 class GlobalTaskLockboxTest
derbyConnector,
new TestDruidLeaderSelector(),
NoopSegmentMetadataCache.instance(),
+ new SegmentsMetadataManagerConfig(null, null, null),
NoopServiceEmitter.instance()
),
objectMapper,
@@ -487,6 +492,7 @@ public class GlobalTaskLockboxTest
derbyConnector,
new TestDruidLeaderSelector(),
NoopSegmentMetadataCache.instance(),
+ new SegmentsMetadataManagerConfig(null, null, null),
NoopServiceEmitter.instance()
),
loadedMapper,
@@ -1308,6 +1314,31 @@ public class GlobalTaskLockboxTest
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<String, List<Interval>> conflictingIntervals =
+
lockbox.getLockedIntervals(ImmutableList.of(requestForExclusiveLowerPriorityLock));
+ Assert.assertTrue(conflictingIntervals.isEmpty());
+ }
+
@Test
public void testGetActiveLocks()
@@ -1848,6 +1879,199 @@ public class GlobalTaskLockboxTest
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<TaskLock> 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<TaskLock> 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 class GlobalTaskLockboxTest
);
}
+ 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 class GlobalTaskLockboxTest
{
final Set<TaskLock> allLocks = getAllLocks();
final Set<TaskLock> 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 class GlobalTaskLockboxTest
{
final Set<TaskLock> allLocks = getAllLocks();
final Set<TaskLock> 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 @@ public class GlobalTaskLockboxTest
{
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 @@ public class GlobalTaskLockboxTest
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<TaskLock> 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 3515790e4a2..2ddd004e583 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.concurrent.Execs;
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 class TaskLockBoxConcurrencyTest
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 190d78bb33c..5432326fd3f 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.io.Closer;
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 class TaskQueueScaleTest
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 1c8d885c924..bed08b4a85d 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.common.task.TaskMetrics;
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 class UnusedSegmentsKillerTest
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 @@ public class UnusedSegmentsKillerTest
SegmentMetadataCache.UsageMode.ALWAYS,
killerConfig
),
- new DefaultTaskConfig(),
taskActionTestKit::createTaskActionClient,
storageCoordinator,
leaderSelector,
@@ -213,7 +211,7 @@ public class UnusedSegmentsKillerTest
@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 class UnusedSegmentsKillerTest
}
@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 class UnusedSegmentsKillerTest
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 @@ public class UnusedSegmentsKillerTest
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 366801ffd7c..8a75d17a80e 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.MonitorScheduler;
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 @@ public abstract class SeekableStreamIndexTaskTestBase
extends EasyMockSupport
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 49e18862c05..ba739e44005 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.core.JsonProcessingException;
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 class IndexerSQLMetadataStorageCoordinator
implements IndexerMetadataStor
@Nullable DateTime maxUsedStatusLastUpdatedTime
)
{
- final List<DataSegment> matchingSegments = inReadOnlyDatasourceTransaction(
- dataSource,
- transaction -> transaction.noCacheSql().findUnusedSegments(
+ final List<DataSegment> matchingSegments = inReadOnlyTransaction(
+ sql -> sql.findUnusedSegments(
dataSource,
interval,
versions,
@@ -2869,50 +2867,18 @@ public class IndexerSQLMetadataStorageCoordinator
implements IndexerMetadataStor
}
/**
- * Performs a read-only transaction using the {@link
SqlSegmentsMetadataQuery},
- * which queries the metadata store directly.
+ * @see
SegmentMetadataTransactionFactory#inReadOnlyNoCacheTransaction(Function)
*/
private <T> T inReadOnlyTransaction(Function<SqlSegmentsMetadataQuery, T>
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> T inWriteTransaction(Function<SqlSegmentsMetadataQuery, T>
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 2be87f57b01..1ca0dfbe948 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 com.google.common.collect.ImmutableSet;
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 abstract class SQLMetadataConnector implements
MetadataStorageConnector
);
}
catch (Exception e) {
- Throwables.propagateIfPossible(e);
+ throwIfUnchecked(e);
throw new RuntimeException(e);
}
}
@@ -193,7 +194,7 @@ public abstract class SQLMetadataConnector implements
MetadataStorageConnector
);
}
catch (Exception e) {
- Throwables.propagateIfPossible(e);
+ throwIfUnchecked(e);
throw new RuntimeException(e);
}
}
@@ -974,7 +975,7 @@ public abstract class SQLMetadataConnector implements
MetadataStorageConnector
);
}
catch (Exception e) {
- Throwables.throwIfUnchecked(e);
+ throwIfUnchecked(e);
throw new RuntimeException(e);
}
}
@@ -1389,6 +1390,15 @@ public abstract class SQLMetadataConnector implements
MetadataStorageConnector
}
}
+ 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 d02bd6bdec2..a2e9e1b4933 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.error.DruidException;
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 d16ba18d0b2..d3f69baa0a0 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.timeline.SegmentTimeline;
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.Map;
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 class SqlSegmentsMetadataQuery
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 class SqlSegmentsMetadataQuery
public boolean markSegmentAsUsed(SegmentId segmentId, DateTime updateTime)
{
+ final List<DataSegmentPlus> 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 @@ public class SqlSegmentsMetadataQuery
DateTime updateTime
)
{
- final List<DataSegment> unusedSegments = new ArrayList<>();
final SegmentTimeline timeline = new SegmentTimeline();
final List<Interval> intervals =
@@ -942,12 +958,13 @@ public class SqlSegmentsMetadataQuery
throw new RuntimeException(e);
}
- try (final CloseableIterator<DataSegment> iterator =
- retrieveUnusedSegments(dataSourceName, intervals, versions, null,
null, null, null)) {
+ final List<DataSegmentPlus> unusedSegments = new ArrayList<>();
+ try (final CloseableIterator<DataSegmentPlus> 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 @@ public class SqlSegmentsMetadataQuery
}
private int markNonOvershadowedSegmentsAsUsed(
- List<DataSegment> unusedSegments,
+ List<DataSegmentPlus> unusedSegments,
SegmentTimeline timeline,
DateTime updateTime
)
{
- Set<SegmentId> nonOvershadowedSegments =
+ final Map<SegmentId, DataSegmentPlus> 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 class SqlSegmentsMetadataQuery
final DateTime updateTime
)
{
- final List<DataSegment> unusedSegments =
retrieveUnusedSegments(dataSource, segmentIds);
+ final List<DataSegmentPlus> unusedSegments =
retrieveUnusedSegmentsPlus(dataSource, segmentIds);
final List<Interval> 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<DataSegment>
usedSegmentsOverlappingUnusedSegmentsIntervals =
retrieveUsedSegments(dataSource, unusedSegmentsIntervals)) {
@@ -997,7 +1016,44 @@ public class SqlSegmentsMetadataQuery
return markNonOvershadowedSegmentsAsUsed(unusedSegments, timeline,
updateTime);
}
- private List<DataSegment> 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<DataSegmentPlus>
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<SegmentId> 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<DataSegmentPlus> retrieveUnusedSegmentsPlus(
final String dataSource,
final Set<SegmentId> segmentIds
)
@@ -1005,12 +1061,11 @@ public class SqlSegmentsMetadataQuery
final List<DataSegmentPlus> retrievedSegments =
retrieveSegmentsById(dataSource, segmentIds);
final Set<SegmentId> unknownSegmentIds = new HashSet<>(segmentIds);
- final List<DataSegment> unusedSegments = new ArrayList<>();
+ final List<DataSegmentPlus> 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 @@ public class SqlSegmentsMetadataQuery
@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 b2deed50438..98fb77ab932 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.JsonCreator;
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 class 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 class 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 924a8c3c54b..af835537ff1 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 e048e10970c..2daa9f906a1 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 @@ public interface SegmentMetadataTransactionFactory
String dataSource,
SegmentMetadataTransaction.Callback<T> callback
);
+
+
+ /**
+ * Performs a read-only transaction which queries the metadata store.
+ */
+ <T> T inReadOnlyNoCacheTransaction(
+ Function<SqlSegmentsMetadataQuery, T> sqlQuery
+ );
+
+ /**
+ * Performs a write transaction which updates the metadata store directly,
+ * and does not affect the cache.
+ */
+ <T> T inReadWriteNoCacheTransaction(
+ Function<SqlSegmentsMetadataQuery, T> 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 0d0d1e42a1e..40c8e6ef6a2 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 com.google.inject.Inject;
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 class SqlSegmentMetadataReadOnlyTransactionFactory
implements SegmentMeta
throw DruidException.defensive("Only Overlord can perform write
transactions on segment metadata.");
}
+ @Override
+ public <T> T inReadOnlyNoCacheTransaction(Function<SqlSegmentsMetadataQuery,
T> sqlQuery)
+ {
+ return connector.retryReadOnlyTransaction(
+ (handle, status) ->
sqlQuery.apply(createSqlQueryForTransaction(handle)),
+ getQuietRetries(),
+ getMaxRetries()
+ );
+ }
+
+ @Override
+ public <T> T
inReadWriteNoCacheTransaction(Function<SqlSegmentsMetadataQuery, T> sqlUpdate)
+ {
+ throw DruidException.defensive("Only Overlord can perform write
transactions on segment metadata.");
+ }
+
protected SegmentMetadataTransaction createSqlTransaction(
String dataSource,
Handle handle,
@@ -98,7 +121,9 @@ public class SqlSegmentMetadataReadOnlyTransactionFactory
implements SegmentMeta
{
return new SqlSegmentMetadataTransaction(
dataSource,
- handle, transactionStatus, connector, tablesConfig, jsonMapper
+ handle,
+ createSqlQueryForTransaction(handle),
+ transactionStatus, connector, tablesConfig, jsonMapper
);
}
@@ -111,4 +136,9 @@ public class SqlSegmentMetadataReadOnlyTransactionFactory
implements SegmentMeta
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 4e144433260..e1ae6f2a957 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 26083e44b67..38d81aa7ef7 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.ServiceEmitter;
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 class SqlSegmentMetadataTransactionFactory extends
SqlSegmentMetadataRead
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 class SqlSegmentMetadataTransactionFactory extends
SqlSegmentMetadataRead
);
}
+ @Override
+ public <T> T
inReadWriteNoCacheTransaction(Function<SqlSegmentsMetadataQuery, T> sqlUpdate)
+ {
+ return connector.retryTransaction(
+ (handle, status) ->
sqlUpdate.apply(createSqlQueryForTransaction(handle)),
+ getQuietRetries(),
+ getMaxRetries()
+ );
+ }
+
private <T> T executeWriteAndClose(
SegmentMetadataTransaction transaction,
SegmentMetadataTransaction.Callback<T> 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 d9d2cfce159..6b62e608521 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 @@ public class HeapMemorySegmentMetadataCache implements
SegmentMetadataCache
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 class HeapMemorySegmentMetadataCache implements
SegmentMetadataCache
)
{
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 @@ public class HeapMemorySegmentMetadataCache implements
SegmentMetadataCache
return inReadOnlyTransaction(
(handle, status) -> sqlFunction.apply(
SqlSegmentsMetadataQuery
- .forHandle(handle, connector, tablesConfig, jsonMapper)
+ .forHandle(handle, connector, tablesConfig, managerConfig,
jsonMapper)
)
);
}
@@ -780,7 +782,7 @@ public class HeapMemorySegmentMetadataCache implements
SegmentMetadataCache
try (
CloseableIterator<DataSegmentPlus> 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 dadad5e4f1b..c807596d882 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 class IndexerSQLMetadataStorageCoordinatorMarkUsedTest
extends IndexerSql
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 d68da314ad1..0ad44376819 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 @@ public class
IndexerSQLMetadataStorageCoordinatorReadOnlyTest extends IndexerSql
transactionFactory = new SqlSegmentMetadataReadOnlyTransactionFactory(
mapper,
derbyConnectorRule.metadataTablesConfigSupplier().get(),
+ new SegmentsMetadataManagerConfig(null, null, null),
derbyConnector
);
} else {
@@ -171,6 +172,7 @@ public class
IndexerSQLMetadataStorageCoordinatorReadOnlyTest extends IndexerSql
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 45f02c3713a..a6a4f8132c6 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.ImmutableSet;
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.Rule;
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 class IndexerSQLMetadataStorageCoordinatorTest
extends IndexerSqlMetadata
derbyConnector,
leaderSelector,
segmentMetadataCache,
+ new SegmentsMetadataManagerConfig(null, cacheMode, null),
emitter
)
{
@@ -599,9 +600,9 @@ public class IndexerSQLMetadataStorageCoordinatorTest
extends IndexerSqlMetadata
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 class IndexerSQLMetadataStorageCoordinatorTest
extends IndexerSqlMetadata
);
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 class IndexerSQLMetadataStorageCoordinatorTest
extends IndexerSqlMetadata
// 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 ba64c709fb1..ec37be03ac4 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 class
IndexerSqlMetadataStorageCoordinatorSchemaPersistenceTest extends
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 be0120162cf..845be8c4824 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 @@ public class IndexerSqlMetadataStorageCoordinatorTestBase
handle,
derbyConnector,
tablesConfig,
+ new
SegmentsMetadataManagerConfig(null, null, null),
mapper
)
.retrieveUnusedSegments(
@@ -385,10 +386,11 @@ public class IndexerSqlMetadataStorageCoordinatorTestBase
MetadataStorageTablesConfig tablesConfig
)
{
+ final SegmentsMetadataManagerConfig managerConfig = new
SegmentsMetadataManagerConfig(null, null, null);
return derbyConnector.inReadOnlyTransaction(
(handle, status) -> {
try (final CloseableIterator<DataSegmentPlus> 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 bedc41f11e8..f7f52077987 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 class SqlSegmentsMetadataManagerTestBase
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 bd480febbdf..b0e8c61cbec 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 @@ package org.apache.druid.metadata;
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.server.coordinator.CreateDataSegments;
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.HashSet;
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 @@ public class SqlSegmentsMetadataQueryTest
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 @@ public class SqlSegmentsMetadataQueryTest
return connector.inReadOnlyTransaction((handle, status) -> {
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<T> iterator = iterableReader.apply(query)) {
return ImmutableSet.copyOf(iterator);
@@ -336,12 +353,33 @@ public class SqlSegmentsMetadataQueryTest
* Executes an update using a {@link SqlSegmentsMetadataQuery} object.
*/
private <T> T update(Function<SqlSegmentsMetadataQuery, T> 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> T updateWithConfig(
+ Function<SqlSegmentsMetadataQuery, T> 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 @@ public class SqlSegmentsMetadataQueryTest
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 @@ public class SqlSegmentsMetadataQueryTest
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<SqlSegmentsMetadataQuery, DataSegment, Integer>
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 <T> void verifyMarkAsUsedThrowsConflictException(
+ BiFunction<SqlSegmentsMetadataQuery, DataSegment, T> 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 37501cb849c..8fbd48ec28f 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.java.util.common.parsers.CloseableIterator;
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 46fc9c835ac..e9dff3b35c9 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 class MetadataManagerModule implements Module
.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);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]