From c6948a2e0b76c665086a46cbca4e1e625329aa14 Mon Sep 17 00:00:00 2001 From: Kartik Khare Date: Mon, 7 Sep 2026 19:37:29 +0530 Subject: [PATCH] Meter upsert compaction segments skipped because replicas disagree UpsertCompactionTask's generator already compares valid doc counts across replicas and refuses to schedule a segment whose replicas disagree. That skip was a log line and nothing else, so the cluster detected divergence, stopped compacting the segment, and told nobody. Two accounts ran for weeks on wrong data as a result. Add ControllerMeter.UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE and emit it from the generator's consensus check. Only an EQUAL-mode disagreement on the valid doc count is metered. The helper also skips a segment when the CRC does not match ZK, when a server is not GOOD, when the CRC will not parse, and when fewer replicas responded than expected. All four mean "ask again later", so metering them would fire on every segment reload and every rolling restart. That is why the meter sits inside selectValidDocIdsMetadataForConsensus, where the reason is still known, rather than at the generator's null check where it is not. The existing 5-argument selectValidDocIdsMetadataForConsensus stays as an overload that meters nothing, so UpsertCompactMergeTask is untouched. The disagreement log line now names both servers and both counts, so the metric gives you the table and the log gives you the segment. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01HyuTYJvV4ARAXR9wSXdjBi --- .../pinot/common/metrics/ControllerMeter.java | 2 + .../plugin/minion/tasks/MinionTaskUtils.java | 26 ++- .../UpsertCompactionTaskGenerator.java | 17 +- .../UpsertCompactionTaskGeneratorTest.java | 149 ++++++++++++++---- 4 files changed, 151 insertions(+), 43 deletions(-) diff --git a/pinot-common/src/main/java/org/apache/pinot/common/metrics/ControllerMeter.java b/pinot-common/src/main/java/org/apache/pinot/common/metrics/ControllerMeter.java index b32d21566d59..a7622cdb951f 100644 --- a/pinot-common/src/main/java/org/apache/pinot/common/metrics/ControllerMeter.java +++ b/pinot-common/src/main/java/org/apache/pinot/common/metrics/ControllerMeter.java @@ -82,6 +82,8 @@ public enum ControllerMeter implements AbstractMetrics.Meter { AUDIT_REQUEST_PAYLOAD_TRUNCATED("count", true), // Upsert compact merge task metrics UPSERT_COMPACT_MERGE_SEGMENT_SKIPPED_CONSENSUS_FAILURE("UpsertCompactMergeSegmentsSkipped", false), + // Segments UpsertCompactionTask refused to schedule because the replicas reported different valid doc counts + UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE("UpsertCompactionSegmentsSkipped", false), // Query workload propagation metrics QUERY_WORKLOAD_PROPAGATION_COUNT("count", true), QUERY_WORKLOAD_PROPAGATION_ERROR("count", true), diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtils.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtils.java index 24878151a476..6a8246d85c52 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtils.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtils.java @@ -40,6 +40,8 @@ import org.apache.pinot.common.auth.AuthProviderUtils; import org.apache.pinot.common.auth.NullAuthProvider; import org.apache.pinot.common.metadata.segment.SegmentZKMetadata; +import org.apache.pinot.common.metrics.ControllerMeter; +import org.apache.pinot.common.metrics.ControllerMetrics; import org.apache.pinot.common.restlet.resources.ValidDocIdsBitmapResponse; import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo; import org.apache.pinot.common.restlet.resources.ValidDocIdsType; @@ -470,6 +472,18 @@ public static RoaringBitmap getValidDocIdFromServerMatchingCrc(String tableNameW public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(String taskType, SegmentZKMetadata segmentZKMetadata, @Nullable List replicas, int expectedReplicaCount, MinionConstants.ValidDocIdsConsensusMode consensusMode) { + return selectValidDocIdsMetadataForConsensus(taskType, segmentZKMetadata, replicas, expectedReplicaCount, + consensusMode, null, null); + } + + /// Same, and additionally reports an EQUAL-mode disagreement on `controllerMetrics` as + /// [ControllerMeter#UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE]. Only that skip reason is metered: the + /// others mean "ask again later", not that the replicas hold different data. + @Nullable + public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(String taskType, + SegmentZKMetadata segmentZKMetadata, @Nullable List replicas, int expectedReplicaCount, + MinionConstants.ValidDocIdsConsensusMode consensusMode, @Nullable ControllerMetrics controllerMetrics, + @Nullable String tableNameWithType) { String segmentName = segmentZKMetadata.getSegmentName(); if (CollectionUtils.isEmpty(replicas)) { return null; @@ -534,9 +548,15 @@ public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(Stri // keeps the generator cheap - it avoids serializing a bitmap per replica back to the controller. ValidDocIdsMetadataInfo first = usableReplicas.get(0); for (int i = 1; i < usableReplicas.size(); i++) { - if (usableReplicas.get(i).getTotalValidDocs() != first.getTotalValidDocs()) { - LOGGER.warn("Replicas disagree on valid doc count for segment: {}, skipping segment for {}", segmentName, - taskType); + ValidDocIdsMetadataInfo other = usableReplicas.get(i); + if (other.getTotalValidDocs() != first.getTotalValidDocs()) { + LOGGER.warn("Replicas disagree on valid doc count for segment: {}, server {} reports {} valid docs and " + + "server {} reports {}, skipping segment for {}", segmentName, first.getInstanceId(), + first.getTotalValidDocs(), other.getInstanceId(), other.getTotalValidDocs(), taskType); + if (controllerMetrics != null && tableNameWithType != null) { + controllerMetrics.addMeteredTableValue(tableNameWithType, + ControllerMeter.UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE, 1L); + } return null; } } diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java index 45a285588c9b..67ed64295a32 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java @@ -27,11 +27,13 @@ import java.util.Map; import java.util.function.Function; import java.util.stream.Collectors; +import javax.annotation.Nullable; import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.helix.task.TaskState; import org.apache.pinot.common.exception.InvalidConfigException; import org.apache.pinot.common.metadata.segment.SegmentZKMetadata; +import org.apache.pinot.common.metrics.ControllerMetrics; import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo; import org.apache.pinot.common.restlet.resources.ValidDocIdsType; import org.apache.pinot.controller.helix.core.PinotHelixResourceManager; @@ -171,8 +173,9 @@ public List generateTasks(List tableConfigs) { completedSegments.stream().collect(Collectors.toMap(SegmentZKMetadata::getSegmentName, Function.identity())); SegmentSelectionResult segmentSelectionResult = - processValidDocIdsMetadata(taskConfigs, completedSegmentsMap, validDocIdsMetadataList, - validDocIdsMetadataResult.getSegmentToExpectedReplicaCount(), consensusMode); + processValidDocIdsMetadata(tableNameWithType, taskConfigs, completedSegmentsMap, validDocIdsMetadataList, + validDocIdsMetadataResult.getSegmentToExpectedReplicaCount(), consensusMode, + _clusterInfoAccessor.getControllerMetrics()); int skippedSegmentsCount = validDocIdsMetadataList.size() - segmentSelectionResult.getSegmentsForCompaction().size() - segmentSelectionResult.getSegmentsForDeletion().size(); @@ -217,10 +220,11 @@ public List generateTasks(List tableConfigs) { } @VisibleForTesting - public static SegmentSelectionResult processValidDocIdsMetadata(Map taskConfigs, - Map completedSegmentsMap, + public static SegmentSelectionResult processValidDocIdsMetadata(String tableNameWithType, + Map taskConfigs, Map completedSegmentsMap, Map> validDocIdsMetadataInfoMap, - Map segmentToReplicaCount, MinionConstants.ValidDocIdsConsensusMode consensusMode) { + Map segmentToReplicaCount, MinionConstants.ValidDocIdsConsensusMode consensusMode, + @Nullable ControllerMetrics controllerMetrics) { double invalidRecordsThresholdPercent = Double.parseDouble( taskConfigs.getOrDefault(UpsertCompactionTask.INVALID_RECORDS_THRESHOLD_PERCENT, String.valueOf(DEFAULT_INVALID_RECORDS_THRESHOLD_PERCENT))); @@ -242,7 +246,8 @@ public static SegmentSelectionResult processValidDocIdsMetadata(Map replicas = validDocIdsMetadataInfoMap.get(segmentName); ValidDocIdsMetadataInfo validDocIdsMetadata = MinionTaskUtils.selectValidDocIdsMetadataForConsensus( MinionConstants.UpsertCompactionTask.TASK_TYPE, segment, replicas, - segmentToReplicaCount.getOrDefault(segmentName, replicas.size()), consensusMode); + segmentToReplicaCount.getOrDefault(segmentName, replicas.size()), consensusMode, controllerMetrics, + tableNameWithType); if (validDocIdsMetadata == null) { continue; } diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java index bc6be0b21eeb..c438b8ddd82d 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java @@ -27,6 +27,8 @@ import java.util.concurrent.TimeUnit; import org.apache.helix.model.IdealState; import org.apache.pinot.common.metadata.segment.SegmentZKMetadata; +import org.apache.pinot.common.metrics.ControllerMeter; +import org.apache.pinot.common.metrics.ControllerMetrics; import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo; import org.apache.pinot.common.restlet.resources.ValidDocIdsType; import org.apache.pinot.common.utils.ServiceStatus; @@ -39,6 +41,7 @@ import org.apache.pinot.spi.config.table.TableType; import org.apache.pinot.spi.config.table.UpsertConfig; import org.apache.pinot.spi.data.Schema; +import org.apache.pinot.spi.metrics.PinotMetricUtils; import org.apache.pinot.spi.utils.CommonConstants; import org.apache.pinot.spi.utils.Enablement; import org.apache.pinot.spi.utils.JsonUtils; @@ -233,14 +236,15 @@ public void testProcessValidDocIdsMetadata() // no completed segments scenario, there shouldn't be any segment selected for compaction UpsertCompactionTaskGenerator.SegmentSelectionResult segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, new HashMap<>(), - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + new HashMap<>(), validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, null); assertEquals(segmentSelectionResult.getSegmentsForCompaction().size(), 0); // test with valid crc and thresholds segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, + null); assertEquals(segmentSelectionResult.getSegmentsForCompaction().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForDeletion().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForCompaction().get(0).getSegmentName(), @@ -250,8 +254,9 @@ public void testProcessValidDocIdsMetadata() // test with a higher invalidRecordsThresholdPercent compactionConfigs = getCompactionConfigs("60", "10"); segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, + null); assertTrue(segmentSelectionResult.getSegmentsForCompaction().isEmpty()); assertEquals(segmentSelectionResult.getSegmentsForDeletion().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForDeletion().get(0), _completedSegment2.getSegmentName()); @@ -259,8 +264,9 @@ public void testProcessValidDocIdsMetadata() // test without an invalidRecordsThresholdPercent compactionConfigs = getCompactionConfigs("0", "10"); segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, + null); assertEquals(segmentSelectionResult.getSegmentsForDeletion().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForCompaction().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForCompaction().get(0).getSegmentName(), @@ -270,8 +276,9 @@ public void testProcessValidDocIdsMetadata() // test without a invalidRecordsThresholdCount compactionConfigs = getCompactionConfigs("30", "0"); segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, + null); assertEquals(segmentSelectionResult.getSegmentsForDeletion().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForCompaction().size(), 1); assertEquals(segmentSelectionResult.getSegmentsForCompaction().get(0).getSegmentName(), @@ -289,8 +296,9 @@ public void testProcessValidDocIdsMetadata() validDocIdsMetadataInfo = JsonUtils.stringToObject(json, new TypeReference<>() { }); segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, + null); // completedSegment is supposed to be filtered out Assert.assertEquals(segmentSelectionResult.getSegmentsForCompaction().size(), 0); @@ -314,8 +322,9 @@ public void testProcessValidDocIdsMetadata() }); compactionConfigs = getCompactionConfigs("30", "0"); segmentSelectionResult = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, validDocIdsMetadataInfo, Map.of(), MinionConstants.ValidDocIdsConsensusMode.UNSAFE, + null); Assert.assertEquals(segmentSelectionResult.getSegmentsForCompaction().size(), 2); Assert.assertEquals(segmentSelectionResult.getSegmentsForDeletion().size(), 0); assertEquals(segmentSelectionResult.getSegmentsForCompaction().get(0).getSegmentName(), @@ -339,51 +348,53 @@ public void testProcessValidDocIdsMetadataConsensus() { meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server2"))); UpsertCompactionTaskGenerator.SegmentSelectionResult result = - UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - equalReplicas, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, equalReplicas, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertEquals(result.getSegmentsForCompaction().size(), 1); Map> unequalReplicas = Map.of(segmentName, List.of( meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), meta(segmentName, 60, 40, 100, crc, ServiceStatus.Status.GOOD, "server2"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - unequalReplicas, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, unequalReplicas, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertTrue(result.getSegmentsForCompaction().isEmpty()); assertTrue(result.getSegmentsForDeletion().isEmpty()); Map> oneResponded = Map.of(segmentName, List.of( meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - oneResponded, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, oneResponded, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertTrue(result.getSegmentsForCompaction().isEmpty()); Map> crcMismatch = Map.of(segmentName, List.of( meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), meta(segmentName, 50, 50, 100, crc + 1, ServiceStatus.Status.GOOD, "server2"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertTrue(result.getSegmentsForCompaction().isEmpty()); Map> unhealthy = Map.of(segmentName, List.of( meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.STARTING, "server2"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - unhealthy, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, unhealthy, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertTrue(result.getSegmentsForCompaction().isEmpty()); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.UNSAFE); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.UNSAFE, null); assertEquals(result.getSegmentsForCompaction().size(), 1); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.MOST_VALID_DOCS); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.MOST_VALID_DOCS, + null); assertTrue(result.getSegmentsForCompaction().isEmpty()); Map> mostValidDocs = Map.of(segmentName, List.of( meta(segmentName, 0, 100, 100, crc, ServiceStatus.Status.GOOD, "server1"), meta(segmentName, 100, 0, 100, crc, ServiceStatus.Status.GOOD, "server2"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, _completedSegmentsMap, - mostValidDocs, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.MOST_VALID_DOCS); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, mostValidDocs, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.MOST_VALID_DOCS, + null); assertTrue(result.getSegmentsForCompaction().isEmpty()); assertTrue(result.getSegmentsForDeletion().isEmpty()); @@ -397,18 +408,88 @@ public void testProcessValidDocIdsMetadataConsensus() { Map> dataCrcMatch = Map.of(segmentName, List.of( metaWithDataCrc(segmentName, 50, 50, 100, 2000, "5000", ServiceStatus.Status.GOOD, "server1"), metaWithDataCrc(segmentName, 50, 50, 100, 2000, "5000", ServiceStatus.Status.GOOD, "server2"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, dataCrcMap, dataCrcMatch, - twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + dataCrcMap, dataCrcMatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertEquals(result.getSegmentsForCompaction().size(), 1); Map> dataCrcMismatch = Map.of(segmentName, List.of( metaWithDataCrc(segmentName, 50, 50, 100, 2000, "9999", ServiceStatus.Status.GOOD, "server1"), metaWithDataCrc(segmentName, 50, 50, 100, 2000, "9999", ServiceStatus.Status.GOOD, "server2"))); - result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(compactionConfigs, dataCrcMap, dataCrcMismatch, - twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL); + result = UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + dataCrcMap, dataCrcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); assertTrue(result.getSegmentsForCompaction().isEmpty()); } + /// Tests that only a genuine replica disagreement raises the divergence meter. A CRC mismatch, an unhealthy + /// server or a short responder list all mean "ask again later", so those skips must leave the meter alone. + @Test + public void testProcessValidDocIdsMetadataConsensusFailureMeter() { + Map compactionConfigs = getCompactionConfigs("1", "10"); + String segmentName = _completedSegment.getSegmentName(); + long crc = _completedSegment.getCrc(); + Map twoReplicas = Map.of(segmentName, 2); + ControllerMetrics controllerMetrics = new ControllerMetrics(PinotMetricUtils.getPinotMetricsRegistry()); + // The registry is shared across tests, so compare deltas rather than absolute counts. + long baseline = meterCount(controllerMetrics); + + Map> agree = Map.of(segmentName, List.of( + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server2"))); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, agree, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, controllerMetrics); + assertEquals(meterCount(controllerMetrics) - baseline, 0, "Agreeing replicas must not raise the meter"); + + Map> crcMismatch = Map.of(segmentName, List.of( + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), + meta(segmentName, 50, 50, 100, crc + 1, ServiceStatus.Status.GOOD, "server2"))); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, crcMismatch, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, + controllerMetrics); + assertEquals(meterCount(controllerMetrics) - baseline, 0, + "A CRC mismatch is a segment reload in flight, not divergence"); + + Map> unhealthy = Map.of(segmentName, List.of( + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.STARTING, "server2"))); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, unhealthy, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, + controllerMetrics); + assertEquals(meterCount(controllerMetrics) - baseline, 0, "A server that is still starting is not divergence"); + + Map> oneResponded = Map.of(segmentName, List.of( + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"))); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, oneResponded, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, + controllerMetrics); + assertEquals(meterCount(controllerMetrics) - baseline, 0, "A short responder list is not divergence"); + + Map> disagree = Map.of(segmentName, List.of( + meta(segmentName, 50, 50, 100, crc, ServiceStatus.Status.GOOD, "server1"), + meta(segmentName, 60, 40, 100, crc, ServiceStatus.Status.GOOD, "server2"))); + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, disagree, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, + controllerMetrics); + assertEquals(meterCount(controllerMetrics) - baseline, 1, + "Replicas reporting different valid doc counts must raise the meter"); + + // MOST_VALID_DOCS picks a winner instead of skipping, so it never reports divergence. + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, disagree, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.MOST_VALID_DOCS, + controllerMetrics); + assertEquals(meterCount(controllerMetrics) - baseline, 1, + "MOST_VALID_DOCS resolves the disagreement, so the meter must not move"); + + // A null ControllerMetrics is the no-metrics path and must not blow up. + UpsertCompactionTaskGenerator.processValidDocIdsMetadata(REALTIME_TABLE_NAME, compactionConfigs, + _completedSegmentsMap, disagree, twoReplicas, MinionConstants.ValidDocIdsConsensusMode.EQUAL, null); + assertEquals(meterCount(controllerMetrics) - baseline, 1, "The null-metrics path must not move the meter"); + } + + private static long meterCount(ControllerMetrics controllerMetrics) { + return controllerMetrics.getMeteredTableValue(REALTIME_TABLE_NAME, + ControllerMeter.UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE).count(); + } + private static ValidDocIdsMetadataInfo meta(String segmentName, long validDocs, long invalidDocs, long totalDocs, long crc, ServiceStatus.Status serverStatus, String instanceId) { return new ValidDocIdsMetadataInfo(segmentName, validDocs, invalidDocs, totalDocs, String.valueOf(crc), null,