From d79b373c70bc14c30c6b7f6792748f187a56bc1d Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Tue, 8 Sep 2026 11:45:08 +0800 Subject: [PATCH 1/7] fix: stabilize task queue and ingestion CI tests --- .../indexing/KafkaBoundedSupervisorTest.java | 25 ++++++++++++-- .../indexing/overlord/TaskQueueScaleTest.java | 20 ++++++----- .../AbstractSegmentMetadataCache.java | 9 +++++ .../schema/BrokerSegmentMetadataCache.java | 11 +++++++ .../BrokerSegmentMetadataCacheTest.java | 33 +++++++++++++++++++ 5 files changed, 87 insertions(+), 11 deletions(-) diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java index 9a38fe5d2e47..d1499b24a62a 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java @@ -29,10 +29,17 @@ import org.apache.druid.query.DruidMetrics; import org.apache.druid.testing.embedded.EmbeddedDruidCluster; import org.apache.druid.testing.embedded.StreamIngestResource; +import org.apache.druid.testing.embedded.tools.EventSerializer; +import org.apache.druid.testing.embedded.tools.JsonEventSerializer; +import org.apache.druid.testing.embedded.tools.StreamGenerator; +import org.apache.druid.testing.embedded.tools.WikipediaStreamEventStreamGenerator; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; /** @@ -186,12 +193,26 @@ private KafkaSupervisorSpec createBoundedKafkaSupervisor( .build(dataSource, topic); } + private void publishRecordsToBothPartitions(String topic) + { + final EventSerializer serializer = new JsonEventSerializer(overlord.bindings().jsonMapper()); + final StreamGenerator generator = new WikipediaStreamEventStreamGenerator(serializer, 100, 100); + final List records = generator.generateEvents(10); + final List> producerRecords = new ArrayList<>(); + // Fixed per-partition end offsets require data in both partitions; Kafka's default partitioner need not balance it. + for (int i = 0; i < records.size(); i++) { + producerRecords.add(new ProducerRecord<>(topic, i % 2, null, records.get(i))); + } + kafkaServer.produceRecordsWithoutTransaction(producerRecords); + Assertions.assertEquals(Map.of("0", 500L, "1", 500L), kafkaServer.getPartitionOffsets(topic)); + } + @Test public void test_boundedSupervisor_withMismatchedMetadata_is_unhealthy() { final String topic = IdUtils.getRandomId(); kafkaServer.createTopicWithPartitions(topic, 2); - publish1kRecords(topic, false); + publishRecordsToBothPartitions(topic); // Get the current end offsets for all partitions Map currentOffsets = kafkaServer.getPartitionOffsets(topic); @@ -265,7 +286,7 @@ public void test_boundedSupervisor_doesNotSilentlyCompleteWhenStaleOffsetExceeds { final String topic = IdUtils.getRandomId(); kafkaServer.createTopicWithPartitions(topic, 2); - publish1kRecords(topic, false); + publishRecordsToBothPartitions(topic); // Run 1: ingest up to offset 100 on each partition and complete. Map startOffsets1 = new HashMap<>(); 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 04670ba28a38..e0c21df78b39 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 @@ -179,12 +179,8 @@ public void doMassLaunchAndExit() throws Exception taskQueue.add(testTask); } - // in theory we can get a race here, since we fetch the counts at separate times + // The running, pending, and waiting counters are separate snapshots and cannot be summed while tasks transition. Assertions.assertEquals(numTasks, taskQueue.getTasks().size(), "all tasks should be known"); - long runningTasks = taskQueue.getRunningTaskCount().values().stream().mapToLong(Long::longValue).sum(); - long pendingTasks = taskQueue.getPendingTaskCount().values().stream().mapToLong(Long::longValue).sum(); - long waitingTasks = taskQueue.getWaitingTaskCount().values().stream().mapToLong(Long::longValue).sum(); - Assertions.assertEquals(numTasks, (runningTasks + pendingTasks + waitingTasks), "all tasks should be known"); // Wait for all tasks to finish. final TaskLookup.CompleteTaskLookup completeTaskLookup = @@ -194,12 +190,18 @@ public void doMassLaunchAndExit() throws Exception Thread.sleep(100); } - Thread.sleep(100); + // Completion is persisted before cleanup finishes. The test timeout bounds this wait. + while (!taskStorage.getActiveTasks().isEmpty() + || taskQueue.getRunningTaskCount().values().stream().anyMatch(count -> count != 0) + || taskQueue.getPendingTaskCount().values().stream().anyMatch(count -> count != 0) + || taskQueue.getWaitingTaskCount().values().stream().anyMatch(count -> count != 0)) { + Thread.sleep(100); + } Assertions.assertEquals(0, taskStorage.getActiveTasks().size(), "no tasks should be active"); - runningTasks = taskQueue.getRunningTaskCount().values().stream().mapToLong(Long::longValue).sum(); - pendingTasks = taskQueue.getPendingTaskCount().values().stream().mapToLong(Long::longValue).sum(); - waitingTasks = taskQueue.getWaitingTaskCount().values().stream().mapToLong(Long::longValue).sum(); + final long runningTasks = taskQueue.getRunningTaskCount().values().stream().mapToLong(Long::longValue).sum(); + final long pendingTasks = taskQueue.getPendingTaskCount().values().stream().mapToLong(Long::longValue).sum(); + final long waitingTasks = taskQueue.getWaitingTaskCount().values().stream().mapToLong(Long::longValue).sum(); Assertions.assertEquals(0, runningTasks, "no tasks should be running"); Assertions.assertEquals(0, pendingTasks, "no tasks should be pending"); Assertions.assertEquals(0, waitingTasks, "no tasks should be waiting"); diff --git a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java index d8677fba9522..48dd50f1b4ea 100644 --- a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java +++ b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java @@ -606,6 +606,7 @@ public void removeSegment(final DataSegment segment) removeSegmentAction(segment.getId()); if (segmentsMap.isEmpty()) { tables.remove(segment.getDataSource()); + removeDataSourceAction(segment.getDataSource()); log.info("dataSource [%s] no longer exists, all metadata removed.", segment.getDataSource()); return null; } else { @@ -625,6 +626,14 @@ public void removeSegment(final DataSegment segment) */ protected abstract void removeSegmentAction(SegmentId segmentId); + /** + * Called under the cache lock after the last segment and its datasource table have been removed. + */ + protected void removeDataSourceAction(String dataSource) + { + // No additional action by default. + } + @VisibleForTesting public void removeServerSegment(final DruidServerMetadata server, final DataSegment segment) { diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java index 26a6c63660d4..4071d6153e00 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java @@ -307,6 +307,17 @@ protected void removeSegmentAction(SegmentId segmentId) // noop, no additional action needed when segment is removed. } + @Override + protected void removeDataSourceAction(String dataSource) + { + // The last-segment callback can remove the table without another schema refresh. + emitMetric( + Metric.DATASOURCE_REMOVED, + 1, + ServiceMetricEvent.builder().setDimension(DruidMetrics.DATASOURCE, dataSource) + ); + } + private Set queryDataSources() { Set dataSources = new HashSet<>(); diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java index 55abaf7b4a39..69c40032cca9 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java @@ -34,6 +34,7 @@ import org.apache.calcite.sql.type.SqlTypeName; import org.apache.druid.client.ImmutableDruidServer; import org.apache.druid.client.InternalQueryConfig; +import org.apache.druid.client.TimelineServerView; import org.apache.druid.client.coordinator.CoordinatorClient; import org.apache.druid.client.coordinator.NoopCoordinatorClient; import org.apache.druid.data.input.InputRow; @@ -720,6 +721,38 @@ public void testNullDatasource() throws IOException, InterruptedException Assertions.assertEquals(5, schema.getSegmentMetadataSnapshot().size()); } + @Test + public void testLastSegmentRemovalEmitsMetricWithoutRefresh() throws IOException + { + final BrokerSegmentMetadataCache schema = new BrokerSegmentMetadataCache( + CalciteTests.createMockQueryLifecycleFactory(walker, conglomerate), + Mockito.mock(TimelineServerView.class), + SEGMENT_CACHE_CONFIG_DEFAULT, + new NoopEscalator(), + new InternalQueryConfig(), + emitter, + new PhysicalDatasourceMetadataFactory(globalTableJoinable, segmentManager), + new NoopCoordinatorClient(), + CentralizedDatasourceSchemaConfig.create() + ); + runningSchema = schema; + final DataSegment segment = walker.getSegments().stream() + .filter(s -> s.getDataSource().equals("foo2")) + .findFirst().orElseThrow(); + schema.addSegment(druidServers.get(0).getMetadata(), segment); + schema.refresh(new HashSet<>(Set.of(segment.getId())), new HashSet<>(Set.of("foo2"))); + Assertions.assertNotNull(schema.getDatasource("foo2")); + emitter.flush(); + + // No background refresh is started: the callback itself must report successful removal. + schema.removeSegment(segment); + Assertions.assertNull(schema.getDatasource("foo2")); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo2"), 1); + + schema.removeSegment(segment); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo2"), 1); + } + @Test public void testAllDatasourcesRebuiltOnDatasourceRemoval() throws IOException, InterruptedException { From ed18502bf6c08777dca5a65dd35cfbb02e2b463c Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Tue, 8 Sep 2026 14:28:56 +0800 Subject: [PATCH 2/7] test: drain ingestion smoke test tasks during cleanup --- .../embedded/indexing/IngestionSmokeTest.java | 54 ++++++++++++++++++- 1 file changed, 52 insertions(+), 2 deletions(-) diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IngestionSmokeTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IngestionSmokeTest.java index d703325ce6e0..b1b9f5d8eac6 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IngestionSmokeTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IngestionSmokeTest.java @@ -24,6 +24,7 @@ import org.apache.druid.common.utils.IdUtils; import org.apache.druid.data.input.impl.CsvInputFormat; import org.apache.druid.data.input.impl.TimestampSpec; +import org.apache.druid.indexer.TaskStatusPlus; import org.apache.druid.indexing.common.task.CompactionTask; import org.apache.druid.indexing.common.task.IndexTask; import org.apache.druid.indexing.common.task.NoopTask; @@ -33,6 +34,9 @@ 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.indexing.overlord.TaskMaster; +import org.apache.druid.indexing.overlord.TaskRunner; +import org.apache.druid.indexing.overlord.TaskRunnerWorkItem; import org.apache.druid.indexing.overlord.supervisor.SupervisorStatus; import org.apache.druid.java.util.common.DateTimes; import org.apache.druid.java.util.common.Intervals; @@ -40,6 +44,7 @@ import org.apache.druid.metadata.storage.postgresql.PostgreSQLMetadataStorageModule; import org.apache.druid.query.DruidMetrics; import org.apache.druid.query.http.SqlTaskStatus; +import org.apache.druid.rpc.indexing.OverlordClient; import org.apache.druid.segment.metadata.Metric; import org.apache.druid.tasklogs.TaskLogStreamer; import org.apache.druid.testing.embedded.EmbeddedBroker; @@ -67,6 +72,7 @@ import java.io.InputStream; import java.nio.charset.StandardCharsets; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -145,7 +151,51 @@ protected EmbeddedDruidCluster addServers(EmbeddedDruidCluster cluster) @AfterEach public void cleanUp() { - markSegmentsAsUnused(dataSource); + try { + final List supervisors = new ArrayList<>(); + cluster.callApi().onLeaderOverlord(OverlordClient::supervisorStatuses).forEachRemaining(supervisors::add); + for (final SupervisorStatus supervisor : supervisors) { + if (dataSource.equals(supervisor.getId())) { + cluster.callApi().onLeaderOverlord(o -> o.terminateSupervisor(supervisor.getId())); + } + } + + cluster.callApi() + .waitForResult(this::cancelTasksForCurrentTest, Set::isEmpty) + .withTimeoutMillis(60_000) + .go(); + } + finally { + markSegmentsAsUnused(dataSource); + } + } + + private Set cancelTasksForCurrentTest() + { + final Set taskIds = new HashSet<>(); + for (final String state : List.of("running", "pending", "waiting")) { + for (final TaskStatusPlus task : cluster.callApi().getTasks(dataSource, state)) { + taskIds.add(task.getId()); + } + } + + // A task can be complete in storage while still occupying a worker slot. + // Docker subclasses may use an external leader, so its runner is not available here. + final Optional taskRunner = overlord.bindings().getInstance(TaskMaster.class).getTaskRunner(); + if (taskRunner.isPresent()) { + final List runnerTasks = new ArrayList<>(taskRunner.get().getPendingTasks()); + runnerTasks.addAll(taskRunner.get().getRunningTasks()); + for (final TaskRunnerWorkItem task : runnerTasks) { + if (dataSource.equals(task.getDataSource())) { + taskIds.add(task.getTaskId()); + } + } + } + + for (final String taskId : taskIds) { + cluster.callApi().onLeaderOverlord(o -> o.cancelTask(taskId)); + } + return taskIds; } protected int markSegmentsAsUnused(String dataSource) @@ -363,7 +413,7 @@ public void test_streamLogs_ofCancelledTask() throws Exception final String taskId = IdUtils.getRandomId(); final long runDurationMillis = 100_000L; cluster.callApi().onLeaderOverlord( - o -> o.runTask(taskId, new NoopTask(taskId, null, null, runDurationMillis, 0L, null)) + o -> o.runTask(taskId, new NoopTask(taskId, null, dataSource, runDurationMillis, 0L, null)) ); eventCollector.latchableEmitter().waitForEvent( From ee5d92b55d98678655c1bb41819702c50bb60f8a Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Tue, 8 Sep 2026 14:45:56 +0800 Subject: [PATCH 3/7] fix: compile supervisor autoscaling tests --- .../supervisor/SupervisorManagerTest.java | 17 +++++++++-------- 1 file changed, 9 insertions(+), 8 deletions(-) diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java index 7fe9a0ad80b2..0b3caa9dd151 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java @@ -32,6 +32,7 @@ import org.apache.druid.data.input.impl.DimensionsSpec; import org.apache.druid.data.input.impl.TimestampSpec; import org.apache.druid.error.DruidException; +import org.apache.druid.error.DruidExceptionMatcher; import org.apache.druid.error.InvalidInput; import org.apache.druid.indexing.common.TaskLockType; import org.apache.druid.indexing.common.task.Tasks; @@ -654,15 +655,15 @@ public void testSimulateAutoscalingUsesLiveTaskCountAboveConfiguredMaximum() thr ); final Object[] data = (Object[]) result.get("data"); - Assert.assertEquals(200, data.length); - Assert.assertTrue(data[0] instanceof Map); + Assertions.assertEquals(200, data.length); + Assertions.assertTrue(data[0] instanceof Map); final Map firstDataPoint = (Map) data[0]; - Assert.assertTrue(firstDataPoint.containsKey("lag")); - Assert.assertTrue(firstDataPoint.containsKey("taskCount")); + Assertions.assertTrue(firstDataPoint.containsKey("lag")); + Assertions.assertTrue(firstDataPoint.containsKey("taskCount")); for (Object dataPoint : data) { - Assert.assertTrue(dataPoint instanceof Map); + Assertions.assertTrue(dataPoint instanceof Map); final Number taskCount = (Number) ((Map) dataPoint).get("taskCount"); - Assert.assertTrue(taskCount.intValue() >= 1 && taskCount.intValue() <= 10); + Assertions.assertTrue(taskCount.intValue() >= 1 && taskCount.intValue() <= 10); } EasyMock.verify(supervisor); } @@ -684,8 +685,8 @@ public void testSimulateAutoscalingRejectsExplicitTaskCountAboveConfiguredMaximu Pair.of(supervisor, new TestBackfillSupervisorSpec(supervisorId, ingestionSpec)) ); - MatcherAssert.assertThat( - Assert.assertThrows( + DruidExceptionMatcher.assertThat( + Assertions.assertThrows( DruidException.class, () -> manager.simulateAutoscaling( supervisorId, From c88b84cec648b9636c070bae1538fa488d5e94a3 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Tue, 8 Sep 2026 21:14:18 +0800 Subject: [PATCH 4/7] fix: clear stale datasource rebuild state on removal --- .../AbstractSegmentMetadataCache.java | 1 + .../BrokerSegmentMetadataCacheTest.java | 29 ++++++++++--------- 2 files changed, 17 insertions(+), 13 deletions(-) diff --git a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java index 48dd50f1b4ea..8b3e3ba3fdfa 100644 --- a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java +++ b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java @@ -605,6 +605,7 @@ public void removeSegment(final DataSegment segment) } removeSegmentAction(segment.getId()); if (segmentsMap.isEmpty()) { + dataSourcesNeedingRebuild.remove(segment.getDataSource()); tables.remove(segment.getDataSource()); removeDataSourceAction(segment.getDataSource()); log.info("dataSource [%s] no longer exists, all metadata removed.", segment.getDataSource()); diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java index 69c40032cca9..19bf6fab092e 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java @@ -722,7 +722,7 @@ public void testNullDatasource() throws IOException, InterruptedException } @Test - public void testLastSegmentRemovalEmitsMetricWithoutRefresh() throws IOException + public void testLastSegmentRemovalClearsRebuildState() throws IOException { final BrokerSegmentMetadataCache schema = new BrokerSegmentMetadataCache( CalciteTests.createMockQueryLifecycleFactory(walker, conglomerate), @@ -736,21 +736,24 @@ public void testLastSegmentRemovalEmitsMetricWithoutRefresh() throws IOException CentralizedDatasourceSchemaConfig.create() ); runningSchema = schema; - final DataSegment segment = walker.getSegments().stream() - .filter(s -> s.getDataSource().equals("foo2")) - .findFirst().orElseThrow(); - schema.addSegment(druidServers.get(0).getMetadata(), segment); - schema.refresh(new HashSet<>(Set.of(segment.getId())), new HashSet<>(Set.of("foo2"))); - Assertions.assertNotNull(schema.getDatasource("foo2")); + final List segments = ImmutableList.of(segment1, segment2); + segments.forEach(segment -> schema.addSegment(druidServers.get(0).getMetadata(), segment)); + schema.refresh( + segments.stream().map(DataSegment::getId).collect(Collectors.toSet()), + new HashSet<>(Set.of("foo")) + ); + Assertions.assertNotNull(schema.getDatasource("foo")); emitter.flush(); - // No background refresh is started: the callback itself must report successful removal. - schema.removeSegment(segment); - Assertions.assertNull(schema.getDatasource("foo2")); - emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo2"), 1); + // Removing a non-last segment marks the datasource for rebuild before the last segment is removed. + schema.removeSegment(segments.get(0)); + schema.removeSegment(segments.get(1)); + Assertions.assertNull(schema.getDatasource("foo")); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo"), 1); - schema.removeSegment(segment); - emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo2"), 1); + // A later refresh must not process stale rebuild state and emit the removal metric again. + schema.refresh(new HashSet<>(), new HashSet<>()); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo"), 1); } @Test From be79adde910381a74c5616bd440eed01e88a689c Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Tue, 8 Sep 2026 22:50:59 +0800 Subject: [PATCH 5/7] fix: satisfy guarded datasource rebuild cleanup --- .../segment/metadata/AbstractSegmentMetadataCache.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java index 8b3e3ba3fdfa..eb3ea0ecdc3b 100644 --- a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java +++ b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java @@ -591,7 +591,7 @@ public void removeSegment(final DataSegment segment) segmentsNeedingRefresh.remove(segment.getId()); unmarkSegmentAsMutable(segment.getId()); - segmentMetadataInfo.compute( + final ConcurrentSkipListMap remainingSegments = segmentMetadataInfo.compute( segment.getDataSource(), (dataSource, segmentsMap) -> { if (segmentsMap == null) { @@ -605,7 +605,6 @@ public void removeSegment(final DataSegment segment) } removeSegmentAction(segment.getId()); if (segmentsMap.isEmpty()) { - dataSourcesNeedingRebuild.remove(segment.getDataSource()); tables.remove(segment.getDataSource()); removeDataSourceAction(segment.getDataSource()); log.info("dataSource [%s] no longer exists, all metadata removed.", segment.getDataSource()); @@ -617,6 +616,9 @@ public void removeSegment(final DataSegment segment) } } ); + if (remainingSegments == null) { + dataSourcesNeedingRebuild.remove(segment.getDataSource()); + } lock.notifyAll(); } From 64694425a909d9b1750a85e70799ce7055cff497 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Thu, 17 Sep 2026 16:50:47 +0800 Subject: [PATCH 6/7] fix: emit dataSource/removed metric at most once per table removal The last-segment callback and an in-flight refresh can both observe the same datasource removal. Gate the metric on the tables.remove() result in both paths so exactly one segment/schemaCache/dataSource/removed event is emitted, regardless of which path removes the table first. --- .../AbstractSegmentMetadataCache.java | 11 ++++++-- .../schema/BrokerSegmentMetadataCache.java | 15 ++++++---- .../BrokerSegmentMetadataCacheTest.java | 28 +++++++++++++++++++ 3 files changed, 46 insertions(+), 8 deletions(-) diff --git a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java index eb3ea0ecdc3b..34d1e755e2f6 100644 --- a/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java +++ b/server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java @@ -605,8 +605,11 @@ public void removeSegment(final DataSegment segment) } removeSegmentAction(segment.getId()); if (segmentsMap.isEmpty()) { - tables.remove(segment.getDataSource()); - removeDataSourceAction(segment.getDataSource()); + // Emit the removal action only if this call actually removed the table, so that a concurrent + // refresh which also finds the datasource gone cannot report the same removal twice. + if (tables.remove(segment.getDataSource()) != null) { + removeDataSourceAction(segment.getDataSource()); + } log.info("dataSource [%s] no longer exists, all metadata removed.", segment.getDataSource()); return null; } else { @@ -630,7 +633,9 @@ public void removeSegment(final DataSegment segment) protected abstract void removeSegmentAction(SegmentId segmentId); /** - * Called under the cache lock after the last segment and its datasource table have been removed. + * Called under the cache lock after the last segment of a datasource has been removed and its table was actually + * removed from {@link #tables} by that removal. It is not called when no table existed for the datasource, so a + * single datasource removal triggers this action at most once even if a refresh observes the removal concurrently. */ protected void removeDataSourceAction(String dataSource) { diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java index 4071d6153e00..69c10b07d9b2 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java @@ -251,11 +251,11 @@ public void refresh(final Set segmentsToRefresh, final Set da final RowSignature rowSignature = buildDataSourceRowSignature(dataSource); if (rowSignature == null) { log.info("datasource [%s] no longer exists, all metadata removed.", dataSource); - tables.remove(dataSource); - emitMetric( - Metric.DATASOURCE_REMOVED, - 1, - ServiceMetricEvent.builder().setDimension(DruidMetrics.DATASOURCE, dataSource)); + // The last-segment callback may already have removed the table and emitted the metric while this refresh + // was in flight. Only emit if this refresh is the one that actually removed the table. + if (tables.remove(dataSource) != null) { + emitDataSourceRemoved(dataSource); + } continue; } @@ -311,6 +311,11 @@ protected void removeSegmentAction(SegmentId segmentId) protected void removeDataSourceAction(String dataSource) { // The last-segment callback can remove the table without another schema refresh. + emitDataSourceRemoved(dataSource); + } + + private void emitDataSourceRemoved(String dataSource) + { emitMetric( Metric.DATASOURCE_REMOVED, 1, diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java index 19bf6fab092e..f8a1d8a5cc78 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java @@ -754,6 +754,34 @@ public void testLastSegmentRemovalClearsRebuildState() throws IOException // A later refresh must not process stale rebuild state and emit the removal metric again. schema.refresh(new HashSet<>(), new HashSet<>()); emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo"), 1); + + // A refresh that had already captured the datasource for rebuild before the last segment was removed + // (or that is handed the datasource explicitly) finds no table left to remove and must not emit again. + schema.refresh(new HashSet<>(), new HashSet<>(Set.of("foo"))); + Assertions.assertNull(schema.getDatasource("foo")); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo"), 1); + } + + @Test + public void testRefreshOfUnknownDatasourceDoesNotEmitRemovalMetric() throws IOException + { + final BrokerSegmentMetadataCache schema = new BrokerSegmentMetadataCache( + CalciteTests.createMockQueryLifecycleFactory(walker, conglomerate), + Mockito.mock(TimelineServerView.class), + SEGMENT_CACHE_CONFIG_DEFAULT, + new NoopEscalator(), + new InternalQueryConfig(), + emitter, + new PhysicalDatasourceMetadataFactory(globalTableJoinable, segmentManager), + new NoopCoordinatorClient(), + CentralizedDatasourceSchemaConfig.create() + ); + runningSchema = schema; + + // No table was ever built for this datasource, so there is nothing to report as removed. + schema.refresh(new HashSet<>(), new HashSet<>(Set.of("never-existed"))); + Assertions.assertNull(schema.getDatasource("never-existed")); + emitter.verifyNotEmitted(Metric.DATASOURCE_REMOVED); } @Test From a59f5f5d4d0856b568dbfe3d31ebd3135b47107d Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Sat, 19 Sep 2026 21:31:27 +0800 Subject: [PATCH 7/7] fix: guard the empty-signature removal metric as well buildDataSourceRowSignature returns an empty (non-null) signature when a datasource only has unrefreshed segments. That branch of refresh() still emitted dataSource/removed unconditionally, so a last-segment callback racing an in-flight refresh could report the removal twice. Apply the same tables.remove() guard, and add a latch-controlled test for the interleaving plus one for a never-initialised datasource. --- .../schema/BrokerSegmentMetadataCache.java | 11 +- .../BrokerSegmentMetadataCacheTest.java | 101 ++++++++++++++++++ 2 files changed, 106 insertions(+), 6 deletions(-) diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java index 69c10b07d9b2..0baa5294f562 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java @@ -264,12 +264,11 @@ public void refresh(final Set segmentsToRefresh, final Set da // and a new datasource is added log.info("datasource [%s] schema has not been initialized yet, " + "check coordinator logs if this message is persistent.", dataSource); - // this is a harmless call - tables.remove(dataSource); - emitMetric( - Metric.DATASOURCE_REMOVED, - 1, - ServiceMetricEvent.builder().setDimension(DruidMetrics.DATASOURCE, dataSource)); + // Usually there is no table to remove here. If there was one, a concurrent last-segment callback may have + // removed it and emitted the metric already, so only emit if this refresh actually removed the table. + if (tables.remove(dataSource) != null) { + emitDataSourceRemoved(dataSource); + } continue; } diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java index f8a1d8a5cc78..57e10f07a5b5 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java @@ -106,6 +106,8 @@ import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; public class BrokerSegmentMetadataCacheTest extends BrokerSegmentMetadataCacheTestBase @@ -784,6 +786,105 @@ public void testRefreshOfUnknownDatasourceDoesNotEmitRemovalMetric() throws IOEx emitter.verifyNotEmitted(Metric.DATASOURCE_REMOVED); } + @Test + public void testRefreshOfUninitializedDatasourceDoesNotEmitRemovalMetric() throws IOException + { + final BrokerSegmentMetadataCache schema = new BrokerSegmentMetadataCache( + CalciteTests.createMockQueryLifecycleFactory(walker, conglomerate), + Mockito.mock(TimelineServerView.class), + SEGMENT_CACHE_CONFIG_DEFAULT, + new NoopEscalator(), + new InternalQueryConfig(), + emitter, + new PhysicalDatasourceMetadataFactory(globalTableJoinable, segmentManager), + new NoopCoordinatorClient(), + CentralizedDatasourceSchemaConfig.create() + ); + runningSchema = schema; + + // The segment is known but has not been refreshed, so the datasource row signature is empty and no table has + // ever been built. Rebuilding it must not report a removal. + schema.addSegment(druidServers.get(0).getMetadata(), segment1); + schema.refresh(new HashSet<>(), new HashSet<>(Set.of("foo"))); + Assertions.assertNull(schema.getDatasource("foo")); + emitter.verifyNotEmitted(Metric.DATASOURCE_REMOVED); + } + + @Test + public void testLastSegmentRemovedDuringRefreshWithUninitializedSegmentEmitsRemovalOnce() throws Exception + { + final AtomicBoolean blockSignatureBuild = new AtomicBoolean(false); + final CountDownLatch signatureBuilt = new CountDownLatch(1); + final CountDownLatch resumeRefresh = new CountDownLatch(1); + final BrokerSegmentMetadataCache schema = new BrokerSegmentMetadataCache( + CalciteTests.createMockQueryLifecycleFactory(walker, conglomerate), + Mockito.mock(TimelineServerView.class), + SEGMENT_CACHE_CONFIG_DEFAULT, + new NoopEscalator(), + new InternalQueryConfig(), + emitter, + new PhysicalDatasourceMetadataFactory(globalTableJoinable, segmentManager), + new NoopCoordinatorClient(), + CentralizedDatasourceSchemaConfig.create() + ) + { + @Override + public RowSignature buildDataSourceRowSignature(String dataSource) + { + final RowSignature rowSignature = super.buildDataSourceRowSignature(dataSource); + if (blockSignatureBuild.get()) { + // Hold the refresh between computing the (empty) signature and acting on it. + signatureBuilt.countDown(); + try { + Assertions.assertTrue(resumeRefresh.await(WAIT_TIMEOUT_SECS, TimeUnit.SECONDS)); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + } + return rowSignature; + } + }; + runningSchema = schema; + + // Build the table from a refreshed segment. + schema.addSegment(druidServers.get(0).getMetadata(), segment1); + schema.refresh(new HashSet<>(Set.of(segment1.getId())), new HashSet<>(Set.of("foo"))); + Assertions.assertNotNull(schema.getDatasource("foo")); + emitter.flush(); + + // Leave only an unrefreshed segment behind, so the next rebuild sees an empty (non-null) row signature. + schema.addSegment(druidServers.get(0).getMetadata(), segment2); + schema.removeSegment(segment1); + + blockSignatureBuild.set(true); + final AtomicReference refreshError = new AtomicReference<>(); + final Thread refreshThread = new Thread(() -> { + try { + schema.refresh(new HashSet<>(), new HashSet<>(Set.of("foo"))); + } + catch (Throwable t) { + refreshError.set(t); + } + }); + refreshThread.start(); + Assertions.assertTrue(signatureBuilt.await(WAIT_TIMEOUT_SECS, TimeUnit.SECONDS)); + + // The last segment disappears while the refresh is in flight: the callback removes the table and emits once. + schema.removeSegment(segment2); + Assertions.assertNull(schema.getDatasource("foo")); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo"), 1); + + // The refresh resumes, finds the table already gone and must not emit a second removal. + resumeRefresh.countDown(); + refreshThread.join(TimeUnit.SECONDS.toMillis(WAIT_TIMEOUT_SECS)); + Assertions.assertFalse(refreshThread.isAlive(), "refresh did not finish"); + Assertions.assertNull(refreshError.get()); + Assertions.assertNull(schema.getDatasource("foo")); + emitter.verifyEmitted(Metric.DATASOURCE_REMOVED, Map.of(DruidMetrics.DATASOURCE, "foo"), 1); + } + @Test public void testAllDatasourcesRebuiltOnDatasourceRemoval() throws IOException, InterruptedException {