diff --git a/pinot-common/src/main/java/org/apache/pinot/common/metrics/MinionMeter.java b/pinot-common/src/main/java/org/apache/pinot/common/metrics/MinionMeter.java index 9baa4c1ee4b3..6e41bfae5ef3 100644 --- a/pinot-common/src/main/java/org/apache/pinot/common/metrics/MinionMeter.java +++ b/pinot-common/src/main/java/org/apache/pinot/common/metrics/MinionMeter.java @@ -42,7 +42,13 @@ public enum MinionMeter implements AbstractMetrics.Meter { TRANSFORMATION_ERROR_COUNT("rows", false), DROPPED_RECORD_COUNT("rows", false), CORRUPTED_RECORD_COUNT("rows", false), - STAR_TREE_INDEX_BUILD_FAILURES("segments", false); + STAR_TREE_INDEX_BUILD_FAILURES("segments", false), + // Upsert compaction CRC / skip observability (issue #13491 residual hardening) + CRC_SKIP_ZK_CHANGED("segments", false), + CRC_MISMATCH_DEEPSTORE("segments", false), + CRC_MISMATCH_SERVER_BITMAP("segments", false), + VALID_DOC_IDS_UNAVAILABLE("segments", false), + COMPACTION_SKIP_EMPTY_VALID_DOCS("segments", false); private final String _meterName; private final String _unit; diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutor.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutor.java index 8a2e9442fc15..c433f37adb9f 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutor.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutor.java @@ -89,8 +89,15 @@ public SegmentConversionResult executeTask(PinotTaskConfig pinotTaskConfig) long currentSegmentCrc = getSegmentCrc(tableNameWithType, segmentName); if (Long.parseLong(originalSegmentCrc) != currentSegmentCrc) { - LOGGER.info("Segment CRC does not match, skip the task. Original CRC: {}, current CRC: {}", originalSegmentCrc, + String skipMessage = String.format( + "Skipped: ZK CRC changed since task generation. Original CRC: %s, current CRC: %s", originalSegmentCrc, currentSegmentCrc); + LOGGER.info("Segment CRC does not match, skip the task. Table: {}, segment: {}. {}", tableNameWithType, + segmentName, skipMessage); + _minionMetrics.addMeteredTableValue(tableNameWithType, MinionMeter.CRC_SKIP_ZK_CHANGED, 1L); + if (_eventObserver != null) { + _eventObserver.notifyProgress(pinotTaskConfig, skipMessage); + } return new SegmentConversionResult.Builder().setTableNameWithType(tableNameWithType).setSegmentName(segmentName) .build(); } 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 3667d7315705..d8f1e02bc731 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 @@ -18,7 +18,6 @@ */ package org.apache.pinot.plugin.minion.tasks; -import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; import java.net.URI; import java.text.SimpleDateFormat; @@ -554,22 +553,30 @@ public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(Stri /// data CRC (a checksum over only the forward index and dictionary, so index/metadata-only changes don't affect /// it) when both copies report one (`>= 0`). A negative data CRC means "not reported". Mirrors the logic of /// `BaseTableDataManager.hasSameCRC` and is used by both the generator's pre-scheduling check and the executor's - /// per-server check. - @VisibleForTesting - static boolean crcMatches(long segmentCrc, long dataCrc, long otherSegmentCrc, long otherDataCrc) { + /// per-server / deepstore checks. + public static boolean crcMatches(long segmentCrc, long dataCrc, long otherSegmentCrc, long otherDataCrc) { if (segmentCrc == otherSegmentCrc) { return true; } return dataCrc >= 0 && otherDataCrc >= 0 && dataCrc == otherDataCrc; } - /// Parses a CRC string, returning `-1` ("unavailable") when it is null or unparseable. - private static long parseCrc(@Nullable String crc) { + /// String overload of {@link #crcMatches(long, long, long, long)}; null/unparseable CRC values are treated as + /// unavailable (`-1`). + public static boolean crcMatches(@Nullable String segmentCrc, @Nullable String dataCrc, + @Nullable String otherSegmentCrc, @Nullable String otherDataCrc) { + return crcMatches(parseCrc(segmentCrc), parseCrc(dataCrc), parseCrc(otherSegmentCrc), parseCrc(otherDataCrc)); + } + + /// Parses a CRC string, returning `-1` ("unavailable") when it is null or unparseable. Values that parse but are + /// negative (e.g. {@link Long#MIN_VALUE} from absent on-disk data CRC) are also treated as unavailable. + public static long parseCrc(@Nullable String crc) { if (crc == null) { return -1; } try { - return Long.parseLong(crc); + long parsed = Long.parseLong(crc); + return parsed >= 0 ? parsed : -1; } catch (NumberFormatException e) { return -1; } diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutor.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutor.java index 2d04d1edaec7..e07a087a6742 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutor.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutor.java @@ -18,9 +18,12 @@ */ package org.apache.pinot.plugin.minion.tasks.upsertcompaction; +import com.google.common.annotations.VisibleForTesting; import java.io.File; import java.util.Map; import org.apache.commons.io.FileUtils; +import org.apache.pinot.common.metadata.ZKMetadataProvider; +import org.apache.pinot.common.metadata.segment.SegmentZKMetadata; import org.apache.pinot.common.metadata.segment.SegmentZKMetadataCustomMapModifier; import org.apache.pinot.common.metrics.MinionMeter; import org.apache.pinot.core.common.MinionConstants; @@ -45,6 +48,19 @@ public class UpsertCompactionTaskExecutor extends BaseSingleSegmentConversionExecutor { private static final Logger LOGGER = LoggerFactory.getLogger(UpsertCompactionTaskExecutor.class); + /** Bounded retries for Check C bitmap fetch to ride out server reload races without re-downloading the segment. */ + @VisibleForTesting + static final int DEFAULT_VALID_DOC_IDS_FETCH_MAX_ATTEMPTS = 3; + + @VisibleForTesting + static final long DEFAULT_VALID_DOC_IDS_FETCH_RETRY_DELAY_MS = 500L; + + @VisibleForTesting + int _validDocIdsFetchMaxAttempts = DEFAULT_VALID_DOC_IDS_FETCH_MAX_ATTEMPTS; + + @VisibleForTesting + long _validDocIdsFetchRetryDelayMs = DEFAULT_VALID_DOC_IDS_FETCH_RETRY_DELAY_MS; + @Override protected SegmentConversionResult convert(PinotTaskConfig pinotTaskConfig, File indexDir, File workingDir) throws Exception { @@ -67,13 +83,8 @@ protected SegmentConversionResult convert(PinotTaskConfig pinotTaskConfig, File String crcFromDeepStorageSegment = segmentMetadata.getCrc(); boolean ignoreCrcMismatch = Boolean.parseBoolean(configs.getOrDefault(UpsertCompactionTask.IGNORE_CRC_MISMATCH_KEY, String.valueOf(UpsertCompactionTask.DEFAULT_IGNORE_CRC_MISMATCH))); - if (!ignoreCrcMismatch && !originalSegmentCrcFromTaskGenerator.equals(crcFromDeepStorageSegment)) { - String message = "Crc mismatched between ZK and deepstore copy of segment: " + segmentName - + ". Expected crc from ZK: " + originalSegmentCrcFromTaskGenerator + ", crc from deepstore: " - + crcFromDeepStorageSegment; - LOGGER.error(message); - throw new IllegalStateException(message); - } + validateDeepStoreCrc(tableNameWithType, segmentName, originalSegmentCrcFromTaskGenerator, crcFromDeepStorageSegment, + segmentMetadata.getDataCrc(), ignoreCrcMismatch); // Executor-only: read comparison mode string from task config (no auth resolution or URL hits). Map taskConfigs = @@ -83,19 +94,15 @@ protected SegmentConversionResult convert(PinotTaskConfig pinotTaskConfig, File MinionConstants.UpsertCompactionTask.DEFAULT_VALID_DOC_IDS_CONSENSUS_MODE) : MinionConstants.UpsertCompactionTask.DEFAULT_VALID_DOC_IDS_CONSENSUS_MODE; RoaringBitmap validDocIds = - MinionTaskUtils.getValidDocIdFromServerMatchingCrc(tableNameWithType, segmentName, validDocIdsTypeStr, - MINION_CONTEXT, originalSegmentCrcFromTaskGenerator, segmentMetadata.getDataCrc(), consensusMode); - if (validDocIds == null) { - // no valid crc match found or no validDocIds obtained from all servers - // error out the task instead of silently failing so that we can track it via task-error metrics - String message = "No validDocIds found from all servers. They either failed to download or did not match crc from" - + " segment copy obtained from deepstore / servers. Expected crc: " + originalSegmentCrcFromTaskGenerator; - LOGGER.error(message); - throw new IllegalStateException(message); - } + fetchValidDocIdsWithRetry(pinotTaskConfig, tableNameWithType, segmentName, validDocIdsTypeStr, + originalSegmentCrcFromTaskGenerator, segmentMetadata.getDataCrc(), consensusMode); if (validDocIds.isEmpty()) { // prevents empty segment generation - LOGGER.info("validDocIds is empty, skip the task. Table: {}, segment: {}", tableNameWithType, segmentName); + String skipMessage = + String.format("Skipped: validDocIds is empty. Table: %s, segment: %s", tableNameWithType, segmentName); + LOGGER.info(skipMessage); + _minionMetrics.addMeteredTableValue(tableNameWithType, MinionMeter.COMPACTION_SKIP_EMPTY_VALID_DOCS, 1L); + _eventObserver.notifyProgress(pinotTaskConfig, skipMessage); if (indexDir.exists() && !FileUtils.deleteQuietly(indexDir)) { LOGGER.warn("Failed to delete input segment: {}", indexDir.getAbsolutePath()); } @@ -137,6 +144,111 @@ protected SegmentConversionResult convert(PinotTaskConfig pinotTaskConfig, File return result; } + /** + * Check B: task-generation (ZK) segment CRC must match the downloaded deepstore/on-disk copy, with the same + * data-CRC fallback used by server matching ({@link MinionTaskUtils#crcMatches}). + */ + @VisibleForTesting + void validateDeepStoreCrc(String tableNameWithType, String segmentName, String expectedSegmentCrc, + String deepstoreSegmentCrc, String deepstoreDataCrc, boolean ignoreCrcMismatch) { + if (ignoreCrcMismatch) { + return; + } + long zkDataCrc = getZkDataCrc(tableNameWithType, segmentName); + if (MinionTaskUtils.crcMatches(MinionTaskUtils.parseCrc(expectedSegmentCrc), zkDataCrc, + MinionTaskUtils.parseCrc(deepstoreSegmentCrc), MinionTaskUtils.parseCrc(deepstoreDataCrc))) { + return; + } + String message = "Crc mismatched between ZK and deepstore copy of segment: " + segmentName + + ". Expected crc from ZK: " + expectedSegmentCrc + ", crc from deepstore: " + deepstoreSegmentCrc + + ", zkDataCrc: " + zkDataCrc + ", deepstoreDataCrc: " + deepstoreDataCrc; + LOGGER.error(message); + _minionMetrics.addMeteredTableValue(tableNameWithType, MinionMeter.CRC_MISMATCH_DEEPSTORE, 1L); + throw new IllegalStateException(message); + } + + /** + * Check C with a short bounded retry so transient server reload races (segment uploaded while a replica is still + * rebuilding upsert metadata) do not fail the task on the first attempt. Does not re-download the segment. + */ + @VisibleForTesting + RoaringBitmap fetchValidDocIdsWithRetry(PinotTaskConfig pinotTaskConfig, String tableNameWithType, String segmentName, + String validDocIdsTypeStr, String expectedSegmentCrc, String expectedDataCrc, String consensusMode) + throws InterruptedException { + int maxAttempts = Math.max(1, _validDocIdsFetchMaxAttempts); + long delayMs = Math.max(0L, _validDocIdsFetchRetryDelayMs); + IllegalStateException lastFailure = null; + + for (int attempt = 1; attempt <= maxAttempts; attempt++) { + try { + RoaringBitmap validDocIds = + MinionTaskUtils.getValidDocIdFromServerMatchingCrc(tableNameWithType, segmentName, validDocIdsTypeStr, + MINION_CONTEXT, expectedSegmentCrc, expectedDataCrc, consensusMode); + if (validDocIds != null) { + if (attempt > 1) { + LOGGER.info("Obtained validDocIds for segment: {} on attempt {}/{}", segmentName, attempt, maxAttempts); + } + return validDocIds; + } + // All servers skipped (UNSAFE) or returned nothing usable. + lastFailure = new IllegalStateException( + "No validDocIds found from all servers. They either failed to download or did not match crc from" + + " segment copy obtained from deepstore / servers. Expected crc: " + expectedSegmentCrc); + LOGGER.warn("validDocIds unavailable for segment: {} on attempt {}/{}", segmentName, attempt, maxAttempts); + } catch (IllegalStateException e) { + lastFailure = e; + LOGGER.warn("validDocIds fetch failed for segment: {} on attempt {}/{}: {}", segmentName, attempt, maxAttempts, + e.getMessage()); + } + + if (attempt < maxAttempts) { + String progress = String.format( + "Retrying validDocIds fetch for segment: %s (attempt %d/%d) after CRC/server mismatch", segmentName, + attempt + 1, maxAttempts); + _eventObserver.notifyProgress(pinotTaskConfig, progress); + if (delayMs > 0L) { + try { + Thread.sleep(delayMs); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + throw ie; + } + } + } + } + + String message = lastFailure != null ? lastFailure.getMessage() + : "No validDocIds found from all servers. Expected crc: " + expectedSegmentCrc; + LOGGER.error(message); + MinionMeter meter = classifyValidDocIdsFailureMeter(message); + _minionMetrics.addMeteredTableValue(tableNameWithType, meter, 1L); + throw lastFailure != null ? lastFailure : new IllegalStateException(message); + } + + @VisibleForTesting + static MinionMeter classifyValidDocIdsFailureMeter(String message) { + if (message != null && message.contains("CRC mismatch")) { + return MinionMeter.CRC_MISMATCH_SERVER_BITMAP; + } + return MinionMeter.VALID_DOC_IDS_UNAVAILABLE; + } + + /** + * ZK data CRC for Check B data-CRC fallback. Returns -1 when metadata is missing or data CRC is not reported. + */ + @VisibleForTesting + long getZkDataCrc(String tableNameWithType, String segmentName) { + SegmentZKMetadata segmentZKMetadata = + ZKMetadataProvider.getSegmentZKMetadata(MINION_CONTEXT.getHelixPropertyStore(), tableNameWithType, segmentName); + if (segmentZKMetadata == null) { + return -1; + } + // Prefer data CRC whenever ZK reports a non-negative value (completed segments may still carry dataCrc after + // commit even when useDataCrc is unset). Mirrors MinionTaskUtils.crcMatches availability rules. + long dataCrc = segmentZKMetadata.getDataCrc(); + return dataCrc >= 0 ? dataCrc : -1; + } + private static SegmentGeneratorConfig getSegmentGeneratorConfig(File workingDir, TableConfig tableConfig, SegmentMetadataImpl segmentMetadata, String segmentName, Schema schema) { SegmentGeneratorConfig config = new SegmentGeneratorConfig(tableConfig, schema); diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutorTest.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutorTest.java index 333f05cfa0e3..d91ae69fa1c1 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutorTest.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/BaseSingleSegmentConversionExecutorTest.java @@ -26,6 +26,7 @@ import java.util.Map; import org.apache.commons.io.FileUtils; import org.apache.pinot.common.metadata.segment.SegmentZKMetadataCustomMapModifier; +import org.apache.pinot.common.metrics.MinionMeter; import org.apache.pinot.common.metrics.MinionMetrics; import org.apache.pinot.core.common.MinionConstants; import org.apache.pinot.core.minion.PinotTaskConfig; @@ -104,6 +105,33 @@ public void setUp() MinionEventObservers.getInstance().addMinionEventObserver(TASK_ID, MinionTaskTestUtils.getMinionProgressObserver()); } + @Test + public void testExecuteTaskSkipsWhenZkCrcChanged() + throws Exception { + MinionMetrics metrics = Mockito.mock(MinionMetrics.class); + // Swap process-global metrics so the CRC_SKIP_ZK_CHANGED meter is observable. + java.lang.reflect.Field field = MinionMetrics.class.getDeclaredField("MINION_METRICS_INSTANCE"); + field.setAccessible(true); + @SuppressWarnings("unchecked") + java.util.concurrent.atomic.AtomicReference ref = + (java.util.concurrent.atomic.AtomicReference) field.get(null); + MinionMetrics previous = ref.getAndSet(metrics); + try { + TestSingleSegmentConversionExecutor executor = new TestSingleSegmentConversionExecutor() { + @Override + protected long getSegmentCrc(String tableNameWithType, String segmentName) { + return SEGMENT_CRC + 1; + } + }; + SegmentConversionResult result = executor.executeTask(createTaskConfig()); + Assert.assertNull(result.getFile()); + Assert.assertEquals(result.getSegmentName(), SEGMENT_NAME); + Mockito.verify(metrics).addMeteredTableValue(TABLE_NAME_WITH_TYPE, MinionMeter.CRC_SKIP_ZK_CHANGED, 1L); + } finally { + ref.set(previous); + } + } + @Test public void testExecuteTaskRethrowsWhenUploadFails() throws Exception { diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtilsTest.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtilsTest.java index b8518c0dd3df..0a2f7fcb9eaf 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtilsTest.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/MinionTaskUtilsTest.java @@ -377,6 +377,22 @@ public void testCrcMatches() { assertFalse(MinionTaskUtils.crcMatches(1000, 5000, 2000, -1)); assertFalse(MinionTaskUtils.crcMatches(1000, -1, 2000, 5000)); assertFalse(MinionTaskUtils.crcMatches(1000, -1, 2000, -1)); + + // String overload (used by Check B deepstore validation). + assertTrue(MinionTaskUtils.crcMatches("1000", "5000", "1000", "9999")); + assertTrue(MinionTaskUtils.crcMatches("1000", "5000", "2000", "5000")); + assertFalse(MinionTaskUtils.crcMatches("1000", "5000", "2000", "9999")); + assertFalse(MinionTaskUtils.crcMatches("1000", null, "2000", "5000")); + } + + @Test + public void testParseCrc() { + assertEquals(MinionTaskUtils.parseCrc("1000"), 1000L); + assertEquals(MinionTaskUtils.parseCrc(null), -1L); + assertEquals(MinionTaskUtils.parseCrc("not-a-number"), -1L); + // Absent on-disk data CRC is serialized as Long.MIN_VALUE; treat as unavailable. + assertEquals(MinionTaskUtils.parseCrc(String.valueOf(Long.MIN_VALUE)), -1L); + assertEquals(MinionTaskUtils.parseCrc("-1"), -1L); } @Test diff --git a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutorTest.java b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutorTest.java index 741a98fa92cf..0c90f4b6be3f 100644 --- a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutorTest.java +++ b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskExecutorTest.java @@ -18,24 +18,87 @@ */ package org.apache.pinot.plugin.minion.tasks.upsertcompaction; +import java.io.File; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.commons.io.FileUtils; import org.apache.helix.HelixAdmin; import org.apache.helix.HelixManager; import org.apache.helix.model.ExternalView; +import org.apache.pinot.common.metrics.MinionMeter; +import org.apache.pinot.common.metrics.MinionMetrics; +import org.apache.pinot.core.common.MinionConstants; +import org.apache.pinot.core.common.MinionConstants.UpsertCompactionTask; +import org.apache.pinot.core.minion.PinotTaskConfig; import org.apache.pinot.minion.MinionContext; +import org.apache.pinot.minion.event.MinionEventObserver; +import org.apache.pinot.plugin.minion.tasks.MinionTaskTestUtils; import org.apache.pinot.plugin.minion.tasks.MinionTaskUtils; +import org.apache.pinot.plugin.minion.tasks.SegmentConversionResult; +import org.apache.pinot.segment.spi.index.metadata.SegmentMetadataImpl; +import org.apache.pinot.spi.config.table.TableConfig; +import org.apache.pinot.spi.config.table.TableTaskConfig; +import org.apache.pinot.spi.config.table.TableType; +import org.apache.pinot.spi.config.table.UpsertConfig; import org.apache.pinot.spi.utils.CommonConstants.Helix.StateModel.SegmentStateModel; +import org.apache.pinot.spi.utils.Enablement; +import org.apache.pinot.spi.utils.builder.TableConfigBuilder; +import org.mockito.MockedConstruction; +import org.mockito.MockedStatic; import org.mockito.Mockito; +import org.roaringbitmap.RoaringBitmap; import org.testng.Assert; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + public class UpsertCompactionTaskExecutorTest { private static final String REALTIME_TABLE_NAME = "testTable_REALTIME"; private static final String SEGMENT_NAME = "testSegment"; private static final String CLUSTER_NAME = "testCluster"; + private static final String TASK_TYPE = UpsertCompactionTask.TASK_TYPE; + private static final String EXPECTED_CRC = "1000"; + private static final String DATA_CRC = "5000"; + + private MinionMetrics _minionMetrics; + private MinionEventObserver _eventObserver; + private File _tempDir; + + @BeforeMethod + public void setUp() + throws Exception { + _minionMetrics = mock(MinionMetrics.class); + // Force the process-global singleton so BaseTaskExecutor picks up the mock. + // register() only wins when current is NOOP; tests may run after other classes registered. + java.lang.reflect.Field field = MinionMetrics.class.getDeclaredField("MINION_METRICS_INSTANCE"); + field.setAccessible(true); + @SuppressWarnings("unchecked") + java.util.concurrent.atomic.AtomicReference ref = + (java.util.concurrent.atomic.AtomicReference) field.get(null); + ref.set(_minionMetrics); + + _eventObserver = MinionTaskTestUtils.getMinionProgressObserver(); + _tempDir = new File(FileUtils.getTempDirectory(), "UpsertCompactionTaskExecutorTest-" + System.nanoTime()); + Assert.assertTrue(_tempDir.mkdirs()); + } + + @AfterMethod + public void tearDown() + throws Exception { + FileUtils.deleteDirectory(_tempDir); + } @Test public void testGetServers() { @@ -64,4 +127,282 @@ public void testGetServers() { () -> MinionTaskUtils.getServers(SEGMENT_NAME, REALTIME_TABLE_NAME, helixManager.getClusterManagmentTool(), helixManager.getClusterName())); } + + @Test + public void testValidateDeepStoreCrcMatch() { + TestableExecutor executor = newExecutor(); + executor.validateDeepStoreCrc(REALTIME_TABLE_NAME, SEGMENT_NAME, EXPECTED_CRC, EXPECTED_CRC, DATA_CRC, false); + verify(_minionMetrics, never()).addMeteredTableValue(anyString(), eq(MinionMeter.CRC_MISMATCH_DEEPSTORE), + anyLong()); + } + + @Test + public void testValidateDeepStoreCrcMismatchThrowsAndMeters() { + TestableExecutor executor = newExecutor(); + executor._zkDataCrc = -1; + try { + executor.validateDeepStoreCrc(REALTIME_TABLE_NAME, SEGMENT_NAME, EXPECTED_CRC, "9999", DATA_CRC, false); + Assert.fail("expected IllegalStateException"); + } catch (IllegalStateException e) { + Assert.assertTrue(e.getMessage().contains("Crc mismatched")); + } + verify(_minionMetrics).addMeteredTableValue(REALTIME_TABLE_NAME, MinionMeter.CRC_MISMATCH_DEEPSTORE, 1L); + } + + @Test + public void testValidateDeepStoreCrcDataCrcFallback() { + // Segment CRCs differ but both sides report the same data CRC → match (index-only drift). + TestableExecutor executor = newExecutor(); + executor._zkDataCrc = Long.parseLong(DATA_CRC); + executor.validateDeepStoreCrc(REALTIME_TABLE_NAME, SEGMENT_NAME, EXPECTED_CRC, "9999", DATA_CRC, false); + verify(_minionMetrics, never()).addMeteredTableValue(anyString(), eq(MinionMeter.CRC_MISMATCH_DEEPSTORE), + anyLong()); + } + + @Test + public void testValidateDeepStoreCrcIgnoreMismatch() { + TestableExecutor executor = newExecutor(); + executor._zkDataCrc = -1; + executor.validateDeepStoreCrc(REALTIME_TABLE_NAME, SEGMENT_NAME, EXPECTED_CRC, "9999", DATA_CRC, true); + verify(_minionMetrics, never()).addMeteredTableValue(anyString(), eq(MinionMeter.CRC_MISMATCH_DEEPSTORE), + anyLong()); + } + + @Test + public void testFetchValidDocIdsRetrySucceedsAfterTransientCrcMismatch() + throws Exception { + TestableExecutor executor = newExecutor(); + executor._validDocIdsFetchMaxAttempts = 3; + executor._validDocIdsFetchRetryDelayMs = 0L; + RoaringBitmap bitmap = new RoaringBitmap(); + bitmap.add(0, 1, 2); + + AtomicInteger calls = new AtomicInteger(); + try (MockedStatic mocked = Mockito.mockStatic(MinionTaskUtils.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when( + () -> MinionTaskUtils.getValidDocIdFromServerMatchingCrc(anyString(), anyString(), anyString(), any(), + anyString(), anyString(), anyString())) + .thenAnswer(inv -> { + if (calls.getAndIncrement() == 0) { + throw new IllegalStateException("CRC mismatch for segment: " + SEGMENT_NAME); + } + return bitmap; + }); + + RoaringBitmap result = + executor.fetchValidDocIdsWithRetry(createTaskConfig(false), REALTIME_TABLE_NAME, SEGMENT_NAME, "SNAPSHOT", + EXPECTED_CRC, DATA_CRC, "EQUAL"); + Assert.assertEquals(result, bitmap); + Assert.assertEquals(calls.get(), 2); + verify(_minionMetrics, never()).addMeteredTableValue(anyString(), eq(MinionMeter.CRC_MISMATCH_SERVER_BITMAP), + anyLong()); + verify(_minionMetrics, never()).addMeteredTableValue(anyString(), eq(MinionMeter.VALID_DOC_IDS_UNAVAILABLE), + anyLong()); + } + } + + @Test + public void testFetchValidDocIdsRetryExhaustedMetersCrcMismatch() { + TestableExecutor executor = newExecutor(); + executor._validDocIdsFetchMaxAttempts = 2; + executor._validDocIdsFetchRetryDelayMs = 0L; + + try (MockedStatic mocked = Mockito.mockStatic(MinionTaskUtils.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when( + () -> MinionTaskUtils.getValidDocIdFromServerMatchingCrc(anyString(), anyString(), anyString(), any(), + anyString(), anyString(), anyString())) + .thenThrow(new IllegalStateException("CRC mismatch for segment: " + SEGMENT_NAME + ", expected: 1000")); + + try { + executor.fetchValidDocIdsWithRetry(createTaskConfig(false), REALTIME_TABLE_NAME, SEGMENT_NAME, "SNAPSHOT", + EXPECTED_CRC, DATA_CRC, "EQUAL"); + Assert.fail("expected IllegalStateException"); + } catch (IllegalStateException e) { + Assert.assertTrue(e.getMessage().contains("CRC mismatch")); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + verify(_minionMetrics).addMeteredTableValue(REALTIME_TABLE_NAME, MinionMeter.CRC_MISMATCH_SERVER_BITMAP, 1L); + } + } + + @Test + public void testFetchValidDocIdsNullAfterRetriesMetersUnavailable() { + TestableExecutor executor = newExecutor(); + executor._validDocIdsFetchMaxAttempts = 2; + executor._validDocIdsFetchRetryDelayMs = 0L; + + try (MockedStatic mocked = Mockito.mockStatic(MinionTaskUtils.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when( + () -> MinionTaskUtils.getValidDocIdFromServerMatchingCrc(anyString(), anyString(), anyString(), any(), + anyString(), anyString(), anyString())) + .thenReturn(null); + + try { + executor.fetchValidDocIdsWithRetry(createTaskConfig(false), REALTIME_TABLE_NAME, SEGMENT_NAME, "SNAPSHOT", + EXPECTED_CRC, DATA_CRC, "UNSAFE"); + Assert.fail("expected IllegalStateException"); + } catch (IllegalStateException e) { + Assert.assertTrue(e.getMessage().contains("No validDocIds")); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + verify(_minionMetrics).addMeteredTableValue(REALTIME_TABLE_NAME, MinionMeter.VALID_DOC_IDS_UNAVAILABLE, 1L); + } + } + + @Test + public void testClassifyValidDocIdsFailureMeter() { + Assert.assertEquals(UpsertCompactionTaskExecutor.classifyValidDocIdsFailureMeter("CRC mismatch for segment: x"), + MinionMeter.CRC_MISMATCH_SERVER_BITMAP); + Assert.assertEquals(UpsertCompactionTaskExecutor.classifyValidDocIdsFailureMeter("No validDocIds found"), + MinionMeter.VALID_DOC_IDS_UNAVAILABLE); + Assert.assertEquals(UpsertCompactionTaskExecutor.classifyValidDocIdsFailureMeter(null), + MinionMeter.VALID_DOC_IDS_UNAVAILABLE); + } + + @Test + public void testConvertEmptyValidDocIdsSkipsWithMeter() + throws Exception { + File indexDir = new File(_tempDir, "indexDir"); + File workingDir = new File(_tempDir, "workingDir"); + Assert.assertTrue(indexDir.mkdirs()); + Assert.assertTrue(workingDir.mkdirs()); + + TestableExecutor executor = newExecutor(); + executor._validDocIdsFetchMaxAttempts = 1; + executor._validDocIdsFetchRetryDelayMs = 0L; + + try (MockedConstruction ignored = Mockito.mockConstruction(SegmentMetadataImpl.class, + (mock, context) -> { + when(mock.getCrc()).thenReturn(EXPECTED_CRC); + when(mock.getDataCrc()).thenReturn(DATA_CRC); + when(mock.getTotalDocs()).thenReturn(10); + }); + MockedStatic mocked = Mockito.mockStatic(MinionTaskUtils.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when(() -> MinionTaskUtils.getValidDocIdsType(any(), any(), anyString())) + .thenReturn(org.apache.pinot.common.restlet.resources.ValidDocIdsType.SNAPSHOT); + mocked.when( + () -> MinionTaskUtils.getValidDocIdFromServerMatchingCrc(anyString(), anyString(), anyString(), any(), + anyString(), anyString(), anyString())) + .thenReturn(new RoaringBitmap()); + + SegmentConversionResult result = executor.convert(createTaskConfig(false), indexDir, workingDir); + Assert.assertNull(result.getFile()); + Assert.assertEquals(result.getSegmentName(), SEGMENT_NAME); + verify(_minionMetrics).addMeteredTableValue(REALTIME_TABLE_NAME, MinionMeter.COMPACTION_SKIP_EMPTY_VALID_DOCS, + 1L); + } + } + + @Test + public void testConvertDeepStoreCrcMismatchThrows() + throws Exception { + File indexDir = new File(_tempDir, "indexDir2"); + File workingDir = new File(_tempDir, "workingDir2"); + Assert.assertTrue(indexDir.mkdirs()); + Assert.assertTrue(workingDir.mkdirs()); + + TestableExecutor executor = newExecutor(); + executor._zkDataCrc = -1; + + try (MockedConstruction ignored = Mockito.mockConstruction(SegmentMetadataImpl.class, + (mock, context) -> { + when(mock.getCrc()).thenReturn("9999"); + when(mock.getDataCrc()).thenReturn(DATA_CRC); + }); + MockedStatic mocked = Mockito.mockStatic(MinionTaskUtils.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when(() -> MinionTaskUtils.getValidDocIdsType(any(), any(), anyString())) + .thenReturn(org.apache.pinot.common.restlet.resources.ValidDocIdsType.SNAPSHOT); + + try { + executor.convert(createTaskConfig(false), indexDir, workingDir); + Assert.fail("expected IllegalStateException"); + } catch (IllegalStateException e) { + Assert.assertTrue(e.getMessage().contains("Crc mismatched")); + } + verify(_minionMetrics).addMeteredTableValue(REALTIME_TABLE_NAME, MinionMeter.CRC_MISMATCH_DEEPSTORE, 1L); + mocked.verify( + () -> MinionTaskUtils.getValidDocIdFromServerMatchingCrc(anyString(), anyString(), anyString(), any(), + anyString(), anyString(), anyString()), never()); + } + } + + @Test + public void testConvertIgnoreCrcMismatchProceedsToValidDocIdsFetch() + throws Exception { + File indexDir = new File(_tempDir, "indexDir3"); + File workingDir = new File(_tempDir, "workingDir3"); + Assert.assertTrue(indexDir.mkdirs()); + Assert.assertTrue(workingDir.mkdirs()); + + TestableExecutor executor = newExecutor(); + executor._zkDataCrc = -1; + executor._validDocIdsFetchMaxAttempts = 1; + executor._validDocIdsFetchRetryDelayMs = 0L; + + try (MockedConstruction ignored = Mockito.mockConstruction(SegmentMetadataImpl.class, + (mock, context) -> { + when(mock.getCrc()).thenReturn("9999"); + when(mock.getDataCrc()).thenReturn(DATA_CRC); + when(mock.getTotalDocs()).thenReturn(10); + }); + MockedStatic mocked = Mockito.mockStatic(MinionTaskUtils.class, Mockito.CALLS_REAL_METHODS)) { + mocked.when(() -> MinionTaskUtils.getValidDocIdsType(any(), any(), anyString())) + .thenReturn(org.apache.pinot.common.restlet.resources.ValidDocIdsType.SNAPSHOT); + mocked.when( + () -> MinionTaskUtils.getValidDocIdFromServerMatchingCrc(anyString(), anyString(), anyString(), any(), + anyString(), anyString(), anyString())) + .thenReturn(new RoaringBitmap()); + + SegmentConversionResult result = executor.convert(createTaskConfig(true), indexDir, workingDir); + Assert.assertNull(result.getFile()); + verify(_minionMetrics, never()).addMeteredTableValue(anyString(), eq(MinionMeter.CRC_MISMATCH_DEEPSTORE), + anyLong()); + verify(_minionMetrics).addMeteredTableValue(REALTIME_TABLE_NAME, MinionMeter.COMPACTION_SKIP_EMPTY_VALID_DOCS, + 1L); + } + } + + private TestableExecutor newExecutor() { + TestableExecutor executor = new TestableExecutor(); + executor.setMinionEventObserver(_eventObserver); + return executor; + } + + private PinotTaskConfig createTaskConfig(boolean ignoreCrcMismatch) { + Map configs = new HashMap<>(); + configs.put(MinionConstants.TABLE_NAME_KEY, REALTIME_TABLE_NAME); + configs.put(MinionConstants.SEGMENT_NAME_KEY, SEGMENT_NAME); + configs.put(MinionConstants.ORIGINAL_SEGMENT_CRC_KEY, EXPECTED_CRC); + configs.put(UpsertCompactionTask.IGNORE_CRC_MISMATCH_KEY, String.valueOf(ignoreCrcMismatch)); + configs.put(UpsertCompactionTask.VALID_DOC_IDS_TYPE, "SNAPSHOT"); + return new PinotTaskConfig(TASK_TYPE, configs); + } + + private static TableConfig createUpsertTableConfig() { + UpsertConfig upsertConfig = new UpsertConfig(UpsertConfig.Mode.FULL); + upsertConfig.setSnapshot(Enablement.ENABLE); + Map> taskTypeConfigs = new HashMap<>(); + taskTypeConfigs.put(TASK_TYPE, Map.of(UpsertCompactionTask.VALID_DOC_IDS_CONSENSUS_MODE_KEY, "EQUAL")); + return new TableConfigBuilder(TableType.REALTIME).setTableName("testTable").setUpsertConfig(upsertConfig) + .setTaskConfig(new TableTaskConfig(taskTypeConfigs)).build(); + } + + /** + * Test double that stubs ZK / table-config lookups so convert() can run without Helix property store. + */ + private static class TestableExecutor extends UpsertCompactionTaskExecutor { + long _zkDataCrc = -1; + + @Override + protected TableConfig getTableConfig(String tableNameWithType) { + return createUpsertTableConfig(); + } + + @Override + long getZkDataCrc(String tableNameWithType, String segmentName) { + return _zkDataCrc; + } + } }