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]