This is an automated email from the ASF dual-hosted git repository.

KKcorps pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git


The following commit(s) were added to refs/heads/master by this push:
     new 747f9470257 Log and count protected upsert revert failures (#19505)
747f9470257 is described below

commit 747f9470257e77ab2e70d1f9f3897adc2d6b5c63
Author: Kartik Khare <[email protected]>
AuthorDate: Thu Oct 1 12:24:21 2026 +0530

    Log and count protected upsert revert failures (#19505)
    
    * Log and count protected upsert revert failures
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Claude-Session: https://claude.ai/code/session_017YJ6Gh5DRvKfpYXniRWai7
    
    * Count keys a concurrent consuming segment took during protected revert
    
    When revert finds the key already owned by another consuming segment, it 
cannot restore the previous
    location. Log it with the UPSERT_METADATA_REVERT_FAILED marker and increment
    UPSERT_METADATA_REVERT_FAILURES, in both ConcurrentMap managers.
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Claude-Session: https://claude.ai/code/session_01LZ977aqwSxDqv41Meb6Adr
    
    ---------
    
    Co-authored-by: Kartik Khare <[email protected]>
    Co-authored-by: Claude Opus 5.5 (1M context) <[email protected]>
---
 .../apache/pinot/common/metrics/ServerMeter.java   |   1 +
 .../upsert/BasePartitionUpsertMetadataManager.java |  13 +-
 ...oncurrentMapPartitionUpsertMetadataManager.java |  21 ++--
 ...nUpsertMetadataManagerForConsistentDeletes.java |  21 ++--
 ...ertMetadataManagerForConsistentDeletesTest.java |  78 ++++++++++++
 ...rrentMapPartitionUpsertMetadataManagerTest.java | 136 +++++++++++++++++++++
 6 files changed, 251 insertions(+), 19 deletions(-)

diff --git 
a/pinot-common/src/main/java/org/apache/pinot/common/metrics/ServerMeter.java 
b/pinot-common/src/main/java/org/apache/pinot/common/metrics/ServerMeter.java
index a02828b392c..0d45a8f1950 100644
--- 
a/pinot-common/src/main/java/org/apache/pinot/common/metrics/ServerMeter.java
+++ 
b/pinot-common/src/main/java/org/apache/pinot/common/metrics/ServerMeter.java
@@ -60,6 +60,7 @@ public enum ServerMeter implements AbstractMetrics.Meter {
   REALTIME_PARTITION_MISMATCH("mismatch", false),
   REALTIME_DEDUP_DROPPED("rows", false),
   DEDUP_PRELOAD_FAILURE("count", false),
+  UPSERT_METADATA_REVERT_FAILURES("failures", false),
   UPSERT_KEYS_IN_WRONG_SEGMENT("rows", false),
   PARTIAL_UPSERT_OUT_OF_ORDER("rows", false),
   PARTIAL_UPSERT_KEYS_NOT_REPLACED("rows", false),
diff --git 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
index 2655bcbad22..cc2e68d0fb7 100644
--- 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
+++ 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
@@ -729,7 +729,18 @@ public abstract class BasePartitionUpsertMetadataManager 
implements PartitionUps
     _logger.info("Inconsistencies noticed for the segment: {} across servers, 
reverting the metadata to resolve...",
         segmentName);
     // Revert the keys in the segment to previous location and remove the 
newly added keys
-    removeSegment(oldSegment, validDocIdsForOldSegment);
+    try {
+      removeSegment(oldSegment, validDocIdsForOldSegment);
+    } catch (RuntimeException e) {
+      String message = "UPSERT_METADATA_REVERT_FAILED: table=" + 
_tableNameWithType + ", partition=" + _partitionId
+          + ", segment=" + segmentName + ". Protected metadata revert did not 
complete; "
+          + "manual reconstruction and replay are required.";
+      _logger.error(message, e);
+      _serverMetrics.addMeteredTableValue(_tableNameWithType, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
+      // Moving the segment to ERROR does not repair partially reverted 
metadata. Report the failure for alerting
+      // instead, so operators can reconstruct and replay the affected 
partition from the failed segment's sequence.
+      return;
+    }
     if (!hasPrevKeyToRecordLocations()) {
       _logger.info("Successfully resolved inconsistency for segment: {} across 
servers", segmentName);
       return;
diff --git 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManager.java
 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManager.java
index f10ce13eead..5860901475b 100644
--- 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManager.java
+++ 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManager.java
@@ -250,26 +250,29 @@ public class ConcurrentMapPartitionUpsertMetadataManager 
extends BasePartitionUp
                       prevDocId, recordInfo);
                   return prevLocation;
                 } catch (Exception e) {
-                  _logger.error("Failed to revert to previous segment: {}, 
removing key", prevSegment.getSegmentName(),
-                      e);
+                  _logger.error("UPSERT_METADATA_REVERT_FAILED: segment={}. 
Failed to revert to previous segment: {}, "
+                      + "removing key", segment.getSegmentName(), 
prevSegment.getSegmentName(), e);
+                  _serverMetrics.addMeteredTableValue(_tableNameWithType,
+                      ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
                   return null;
                 }
               } else {
                 // Should not happen
-                _logger.error("Failed to find valid doc ids in previous 
segment: {}, removing key",
-                    prevSegment.getSegmentName());
+                _logger.error("UPSERT_METADATA_REVERT_FAILED: segment={}. 
Failed to find valid doc ids in previous "
+                    + "segment: {}, removing key", segment.getSegmentName(), 
prevSegment.getSegmentName());
+                _serverMetrics.addMeteredTableValue(_tableNameWithType, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
                 return null;
               }
             } else if (recordLocation.getSegment() instanceof 
ImmutableSegmentImpl) {
               // The consuming segment's key is in a different immutable 
segment
               _previousKeyToRecordLocationMap.remove(pk);
             } else {
-              _logger.warn(
-                  "Consuming segment: {} has added the primary key for docId: 
{} from the segment: {}, suggesting"
-                      + " that consumption is occurring concurrently with 
segment replacement, which is undesirable "
-                      + "for consistency between replicas for the table: {}.",
-                  recordLocation.getSegment().getSegmentName(), 
primaryKeyEntry.getKey(), segment.getSegmentName(),
+              _logger.warn("UPSERT_METADATA_REVERT_FAILED: segment={}. 
Consuming segment: {} has added the primary "
+                      + "key for docId: {}, suggesting that consumption is 
occurring concurrently with segment "
+                      + "replacement, which is undesirable for consistency 
between replicas for the table: {}.",
+                  segment.getSegmentName(), 
recordLocation.getSegment().getSegmentName(), primaryKeyEntry.getKey(),
                   _tableNameWithType);
+              _serverMetrics.addMeteredTableValue(_tableNameWithType, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
             }
             return recordLocation;
           });
diff --git 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes.java
 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes.java
index c54dac31056..19c0c301bfc 100644
--- 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes.java
+++ 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes.java
@@ -380,26 +380,29 @@ public class 
ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes
                       prevLocation.getComparisonValue(),
                       
RecordLocation.decrementSegmentCount(prevLocation.getDistinctSegmentCount()));
                 } catch (Exception e) {
-                  _logger.error("Failed to revert to previous segment: {}, 
removing key", prevSegment.getSegmentName(),
-                      e);
+                  _logger.error("UPSERT_METADATA_REVERT_FAILED: segment={}. 
Failed to revert to previous segment: {}, "
+                      + "removing key", segment.getSegmentName(), 
prevSegment.getSegmentName(), e);
+                  _serverMetrics.addMeteredTableValue(_tableNameWithType,
+                      ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
                   return null;
                 }
               } else {
                 // Should not happen
-                _logger.error("Failed to find valid doc ids in previous 
segment: {}, removing key",
-                    prevSegment.getSegmentName());
+                _logger.error("UPSERT_METADATA_REVERT_FAILED: segment={}. 
Failed to find valid doc ids in previous "
+                    + "segment: {}, removing key", segment.getSegmentName(), 
prevSegment.getSegmentName());
+                _serverMetrics.addMeteredTableValue(_tableNameWithType, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
                 return null;
               }
             } else if (recordLocation.getSegment() instanceof 
ImmutableSegmentImpl) {
               // The consuming segment's key is in a different immutable 
segment
               _previousKeyToRecordLocationMap.remove(pk);
             } else {
-              _logger.warn(
-                  "Consuming segment: {} has added the primary key for docId: 
{} from the segment: {}, suggesting"
-                      + " that consumption is occurring concurrently with 
segment replacement, which is undesirable "
-                      + "for consistency between replicas for the table: {}.",
-                  recordLocation.getSegment().getSegmentName(), 
primaryKeyEntry.getKey(), segment.getSegmentName(),
+              _logger.warn("UPSERT_METADATA_REVERT_FAILED: segment={}. 
Consuming segment: {} has added the primary "
+                      + "key for docId: {}, suggesting that consumption is 
occurring concurrently with segment "
+                      + "replacement, which is undesirable for consistency 
between replicas for the table: {}.",
+                  segment.getSegmentName(), 
recordLocation.getSegment().getSegmentName(), primaryKeyEntry.getKey(),
                   _tableNameWithType);
+              _serverMetrics.addMeteredTableValue(_tableNameWithType, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
             }
             if (!uniquePrimaryKeys.add(pk)) {
               return recordLocation;
diff --git 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletesTest.java
 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletesTest.java
index 52a688907f8..5b320aec117 100644
--- 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletesTest.java
+++ 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletesTest.java
@@ -31,6 +31,7 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.atomic.AtomicBoolean;
 import javax.annotation.Nullable;
 import org.apache.commons.io.FileUtils;
+import org.apache.pinot.common.metrics.ServerMeter;
 import org.apache.pinot.common.metrics.ServerMetrics;
 import org.apache.pinot.common.utils.LLCSegmentName;
 import org.apache.pinot.common.utils.UploadedRealtimeSegmentName;
@@ -38,6 +39,7 @@ import 
org.apache.pinot.segment.local.data.manager.TableDataManager;
 import org.apache.pinot.segment.local.indexsegment.immutable.EmptyIndexSegment;
 import 
org.apache.pinot.segment.local.indexsegment.immutable.ImmutableSegmentImpl;
 import org.apache.pinot.segment.local.segment.readers.PinotSegmentColumnReader;
+import 
org.apache.pinot.segment.local.upsert.ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes.RecordLocation;
 import org.apache.pinot.segment.local.utils.HashUtils;
 import org.apache.pinot.segment.spi.ColumnMetadata;
 import org.apache.pinot.segment.spi.IndexSegment;
@@ -57,6 +59,7 @@ import org.apache.pinot.spi.data.Schema;
 import org.apache.pinot.spi.data.readers.PrimaryKey;
 import org.apache.pinot.spi.utils.ByteArray;
 import org.apache.pinot.spi.utils.BytesUtils;
+import org.apache.pinot.spi.utils.ConsumingSegmentConsistencyModeListener;
 import org.apache.pinot.spi.utils.builder.TableNameBuilder;
 import org.apache.pinot.util.TestUtils;
 import org.mockito.MockedConstruction;
@@ -64,6 +67,7 @@ import org.roaringbitmap.buffer.MutableRoaringBitmap;
 import org.testng.annotations.AfterClass;
 import org.testng.annotations.BeforeClass;
 import org.testng.annotations.BeforeMethod;
+import org.testng.annotations.DataProvider;
 import org.testng.annotations.Test;
 
 import static org.mockito.ArgumentMatchers.any;
@@ -71,6 +75,9 @@ import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 import static org.testng.Assert.*;
 
@@ -1271,6 +1278,77 @@ public class 
ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletesTest
         "37fab5ef0ea39711feabcdc623cb8a4e");
   }
 
+  @DataProvider
+  public Object[][] revertFailureCases() {
+    return new Object[][]{
+        {"reader"}, {"bitmap"}, {"concurrent"}, {"none"}
+    };
+  }
+
+  @Test(dataProvider = "revertFailureCases")
+  public void testRevertFailureReporting(String failure) {
+    ConsumingSegmentConsistencyModeListener listener = 
ConsumingSegmentConsistencyModeListener.getInstance();
+    ConsumingSegmentConsistencyModeListener.Mode originalMode = 
listener.getConsistencyMode();
+    ServerMetrics originalMetrics = ServerMetrics.get();
+    ServerMetrics metrics = mock(ServerMetrics.class);
+    ServerMetrics.deregister();
+    ServerMetrics.register(metrics);
+    listener.setMode(ConsumingSegmentConsistencyModeListener.Mode.PROTECTED);
+    try {
+      UpsertContext context = 
_contextBuilder.setDropOutOfOrderRecord(true).setHashFunction(HashFunction.NONE).build();
+      ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes manager =
+          new 
ConcurrentMapPartitionUpsertMetadataManagerForConsistentDeletes(REALTIME_TABLE_NAME,
 0, context);
+      ThreadSafeMutableRoaringBitmap validDocIds = new 
ThreadSafeMutableRoaringBitmap();
+      validDocIds.add(0);
+      MutableSegment segment = mockMutableSegmentWithDataSource(1, 
validDocIds, null, new int[]{10});
+      ImmutableSegmentImpl previousSegment = mock(ImmutableSegmentImpl.class);
+      when(previousSegment.getSegmentName()).thenReturn(getSegmentName(0));
+      ThreadSafeMutableRoaringBitmap previousValidDocIds = new 
ThreadSafeMutableRoaringBitmap();
+      
when(previousSegment.getValidDocIds()).thenReturn(failure.equals("bitmap") ? 
null : previousValidDocIds);
+      PrimaryKey key = makePrimaryKey(10);
+      // "concurrent": the next consuming segment took the key while this 
segment was being replaced
+      MutableSegment nextSegment = mock(MutableSegment.class);
+      when(nextSegment.getSegmentName()).thenReturn(getSegmentName(2));
+      boolean concurrent = failure.equals("concurrent");
+      manager._primaryKeyToRecordLocationMap.put(key, new 
RecordLocation(concurrent ? nextSegment : segment, 0,
+          concurrent ? 300 : 200, 2));
+      manager._previousKeyToRecordLocationMap.put(key, new 
RecordLocation(previousSegment, 0, 100, 2));
+      manager._trackedSegments.add(segment);
+      try (MockedConstruction<UpsertUtils.RecordInfoReader> readers =
+          mockConstruction(UpsertUtils.RecordInfoReader.class, (reader, 
construction) -> {
+            if (failure.equals("reader")) {
+              when(reader.getRecordInfo(0)).thenThrow(new 
IllegalStateException("previous reader failed"));
+            } else {
+              when(reader.getRecordInfo(0)).thenReturn(new RecordInfo(key, 0, 
100, false));
+            }
+          })) {
+        manager.removeSegment(segment);
+        assertFalse(manager._trackedSegments.contains(segment));
+        assertEquals(readers.constructed().size(), failure.equals("bitmap") || 
concurrent ? 0 : 1);
+      }
+      boolean successfulRevert = failure.equals("none");
+      if (successfulRevert) {
+        checkRecordLocation(manager._primaryKeyToRecordLocationMap, 10, 
previousSegment, 0, 100, 1, HashFunction.NONE);
+        assertTrue(previousValidDocIds.contains(0));
+        assertFalse(validDocIds.contains(0));
+      } else if (concurrent) {
+        checkRecordLocation(manager._primaryKeyToRecordLocationMap, 10, 
nextSegment, 0, 300, 1, HashFunction.NONE);
+      } else {
+        assertFalse(manager._primaryKeyToRecordLocationMap.containsKey(key),
+            "Preserve the existing key-removal fallback");
+      }
+      // Removal reverts keys directly, so a key owned by the next consuming 
segment keeps its previous location
+      assertEquals(manager._previousKeyToRecordLocationMap.containsKey(key), 
concurrent);
+      verify(metrics, times(successfulRevert ? 0 : 1))
+          .addMeteredTableValue(REALTIME_TABLE_NAME, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
+      verify(context.getTableDataManager(), 
never()).addSegmentError(anyString(), any());
+    } finally {
+      listener.setMode(originalMode);
+      ServerMetrics.deregister();
+      ServerMetrics.register(originalMetrics);
+    }
+  }
+
   @Test
   public void testRevertOnlyAppliesForConsumingSegmentSeal()
       throws IOException {
diff --git 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
index 06b97a2f1dc..d76c64d34c9 100644
--- 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
+++ 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
@@ -31,6 +31,7 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.atomic.AtomicBoolean;
 import javax.annotation.Nullable;
 import org.apache.commons.io.FileUtils;
+import org.apache.pinot.common.metrics.ServerMeter;
 import org.apache.pinot.common.metrics.ServerMetrics;
 import org.apache.pinot.common.utils.LLCSegmentName;
 import org.apache.pinot.common.utils.UploadedRealtimeSegmentName;
@@ -77,6 +78,7 @@ import org.roaringbitmap.buffer.MutableRoaringBitmap;
 import org.testng.annotations.AfterClass;
 import org.testng.annotations.BeforeClass;
 import org.testng.annotations.BeforeMethod;
+import org.testng.annotations.DataProvider;
 import org.testng.annotations.Test;
 
 import static org.mockito.ArgumentMatchers.any;
@@ -84,8 +86,12 @@ import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 import static org.testng.Assert.*;
@@ -2406,6 +2412,136 @@ public class 
ConcurrentMapPartitionUpsertMetadataManagerTest {
     upsertMetadataManager.close();
   }
 
+  @DataProvider
+  public Object[][] metadataRevertFailureCases() {
+    return new Object[][]{
+        {ConsumingSegmentConsistencyModeListener.Mode.PROTECTED, false},
+        {ConsumingSegmentConsistencyModeListener.Mode.PROTECTED, true},
+        {ConsumingSegmentConsistencyModeListener.Mode.RESTRICTED, false},
+        {ConsumingSegmentConsistencyModeListener.Mode.RESTRICTED, true}
+    };
+  }
+
+  @Test(dataProvider = "metadataRevertFailureCases")
+  public void 
testProtectedRevertFailuresAreReportedWithoutThrowing(ConsumingSegmentConsistencyModeListener.Mode
 mode,
+      boolean replacement) {
+    ConsumingSegmentConsistencyModeListener listener = 
ConsumingSegmentConsistencyModeListener.getInstance();
+    ConsumingSegmentConsistencyModeListener.Mode originalMode = 
listener.getConsistencyMode();
+    ServerMetrics originalMetrics = ServerMetrics.get();
+    ServerMetrics metrics = mock(ServerMetrics.class);
+    ServerMetrics.deregister();
+    ServerMetrics.register(metrics);
+    listener.setMode(mode);
+    try {
+      UpsertContext context = 
_contextBuilder.setDropOutOfOrderRecord(true).build();
+      ConcurrentMapPartitionUpsertMetadataManager manager =
+          spy(new 
ConcurrentMapPartitionUpsertMetadataManager(REALTIME_TABLE_NAME, 0, context));
+      MutableSegment segment = mock(MutableSegment.class);
+      String segmentName = "testTable__0__1__0";
+      when(segment.getSegmentName()).thenReturn(segmentName);
+      ThreadSafeMutableRoaringBitmap validDocIds = new 
ThreadSafeMutableRoaringBitmap();
+      validDocIds.add(0);
+      when(segment.getValidDocIds()).thenReturn(validDocIds);
+      manager._trackedSegments.add(segment);
+      RuntimeException failure = new RuntimeException("removal failed");
+      doThrow(failure).when(manager).removeSegment(eq(segment), 
any(MutableRoaringBitmap.class));
+
+      Runnable removeOrReplaceSegment = () -> {
+        if (replacement) {
+          ImmutableSegmentImpl newSegment = mock(ImmutableSegmentImpl.class);
+          when(newSegment.getSegmentName()).thenReturn(segmentName);
+          manager.replaceSegment(newSegment, null, null, null, segment);
+        } else {
+          manager.removeSegment(segment);
+        }
+      };
+      boolean report = mode == 
ConsumingSegmentConsistencyModeListener.Mode.PROTECTED;
+      if (report) {
+        removeOrReplaceSegment.run();
+        if (!replacement) {
+          assertFalse(manager._trackedSegments.contains(segment), "Failed 
revert must not prevent segment offload");
+        }
+      } else {
+        RuntimeException thrown = expectThrows(RuntimeException.class, 
removeOrReplaceSegment::run);
+        assertSame(thrown, failure, "Preserve the original failure outside 
protected revert");
+      }
+      verify(manager).removeSegment(eq(segment), 
any(MutableRoaringBitmap.class));
+      verify(metrics, times(report ? 1 : 0))
+          .addMeteredTableValue(REALTIME_TABLE_NAME, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
+      verify(context.getTableDataManager(), 
never()).addSegmentError(anyString(), any());
+    } finally {
+      listener.setMode(originalMode);
+      ServerMetrics.deregister();
+      ServerMetrics.register(originalMetrics);
+    }
+  }
+
+  @DataProvider
+  public Object[][] handledRevertFailureCases() {
+    return new Object[][]{
+        {"reader"}, {"bitmap"}, {"concurrent"}
+    };
+  }
+
+  @Test(dataProvider = "handledRevertFailureCases")
+  public void testHandledRevertFailuresKeepExistingBehavior(String failure) {
+    ConsumingSegmentConsistencyModeListener listener = 
ConsumingSegmentConsistencyModeListener.getInstance();
+    ConsumingSegmentConsistencyModeListener.Mode originalMode = 
listener.getConsistencyMode();
+    ServerMetrics originalMetrics = ServerMetrics.get();
+    ServerMetrics metrics = mock(ServerMetrics.class);
+    ServerMetrics.deregister();
+    ServerMetrics.register(metrics);
+    listener.setMode(ConsumingSegmentConsistencyModeListener.Mode.PROTECTED);
+    try {
+      UpsertContext context = 
_contextBuilder.setDropOutOfOrderRecord(true).setHashFunction(HashFunction.NONE).build();
+      ConcurrentMapPartitionUpsertMetadataManager manager =
+          new ConcurrentMapPartitionUpsertMetadataManager(REALTIME_TABLE_NAME, 
0, context);
+      ThreadSafeMutableRoaringBitmap validDocIds = new 
ThreadSafeMutableRoaringBitmap();
+      validDocIds.add(0);
+      MutableSegment segment = mockMutableSegmentWithDataSource(1, 
validDocIds, null, new int[]{10});
+      ImmutableSegmentImpl previousSegment = mock(ImmutableSegmentImpl.class);
+      when(previousSegment.getSegmentName()).thenReturn(getSegmentName(0));
+      when(previousSegment.getValidDocIds()).thenReturn(
+          failure.equals("bitmap") ? null : new 
ThreadSafeMutableRoaringBitmap());
+      PrimaryKey key = makePrimaryKey(10);
+      // "concurrent": the next consuming segment took the key while this 
segment was being replaced
+      MutableSegment nextSegment = mock(MutableSegment.class);
+      when(nextSegment.getSegmentName()).thenReturn(getSegmentName(2));
+      boolean concurrent = failure.equals("concurrent");
+      manager._primaryKeyToRecordLocationMap.put(key,
+          concurrent ? new RecordLocation(nextSegment, 0, 300) : new 
RecordLocation(segment, 0, 200));
+      manager._previousKeyToRecordLocationMap.put(key, new 
RecordLocation(previousSegment, 0, 100));
+      manager._trackedSegments.add(segment);
+      try (MockedConstruction<UpsertUtils.RecordInfoReader> readers =
+          mockConstruction(UpsertUtils.RecordInfoReader.class, (reader, 
construction) -> {
+            if (failure.equals("reader")) {
+              when(reader.getRecordInfo(0)).thenThrow(new 
IllegalStateException("previous reader failed"));
+            } else {
+              when(reader.getRecordInfo(0)).thenReturn(new RecordInfo(key, 0, 
100, false));
+            }
+          })) {
+        // Use the actual backend removal. Existing reader/bitmap fallbacks 
must finish without throwing.
+        manager.removeSegment(segment);
+        assertFalse(manager._trackedSegments.contains(segment));
+        assertEquals(readers.constructed().size(), failure.equals("bitmap") || 
concurrent ? 0 : 1);
+      }
+      if (concurrent) {
+        
assertSame(manager._primaryKeyToRecordLocationMap.get(key).getSegment(), 
nextSegment,
+            "Keep the key in the next consuming segment");
+      } else {
+        assertFalse(manager._primaryKeyToRecordLocationMap.containsKey(key),
+            "Preserve the existing key-removal fallback");
+      }
+      assertTrue(manager._previousKeyToRecordLocationMap.isEmpty());
+      verify(metrics).addMeteredTableValue(REALTIME_TABLE_NAME, 
ServerMeter.UPSERT_METADATA_REVERT_FAILURES, 1);
+      verify(context.getTableDataManager(), 
never()).addSegmentError(anyString(), any());
+    } finally {
+      listener.setMode(originalMode);
+      ServerMetrics.deregister();
+      ServerMetrics.register(originalMetrics);
+    }
+  }
+
   @Test
   public void testProtectedModeRevertsMetadataForConsumingSegmentSeal()
       throws IOException {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to