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,