This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang 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 fe8cd0261ec Change SegmentMetadata.getCrc and getDataCrc to return
long (#19629)
fe8cd0261ec is described below
commit fe8cd0261ece07f59204b60f0780202ce0ebc1fb
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Tue Sep 22 22:48:51 2026 -0700
Change SegmentMetadata.getCrc and getDataCrc to return long (#19629)
---
.../metadata/segment/SegmentZKMetadataUtils.java | 6 +-
.../pinot/controller/api/upload/ZKOperator.java | 2 +-
.../realtime/provisioning/MemoryEstimator.java | 4 +-
.../api/PinotSegmentRestletResourceTest.java | 4 +-
.../controller/api/upload/ZKOperatorTest.java | 20 +-
.../LineageDeleteInterleavingIntegrationTest.java | 2 +-
.../PinotLLCRealtimeSegmentManagerTest.java | 14 +-
.../controller/utils/SegmentMetadataMockUtils.java | 32 +-
.../validation/ValidationManagerTest.java | 2 +-
.../core/data/manager/BaseTableDataManager.java | 17 +-
...ableDataManagerEnqueueSegmentToReplaceTest.java | 2 +-
.../data/manager/BaseTableDataManagerTest.java | 24 +-
.../offline/DimensionTableDataManagerTest.java | 2 +-
.../realtime/RealtimeSegmentDataManagerTest.java | 2 +-
.../tasks/BaseSingleSegmentConversionExecutor.java | 2 +-
.../pinot/plugin/minion/tasks/MinionTaskUtils.java | 8 +-
.../UpsertCompactionTaskExecutor.java | 6 +-
.../UpsertCompactMergeTaskExecutor.java | 11 +-
.../BaseSingleSegmentConversionExecutorTest.java | 4 +-
.../plugin/minion/tasks/MinionTaskUtilsTest.java | 46 +--
.../UpsertCompactMergeTaskExecutorTest.java | 44 +-
.../local/data/manager/SegmentDataManager.java | 2 +-
.../immutable/ImmutableSegmentLoader.java | 2 +-
.../SegmentCrcVirtualColumnProvider.java | 23 +-
.../converter/RealtimeSegmentConverterTest.java | 2 +-
.../segment/index/SegmentMetadataImplTest.java | 2 +-
.../SegmentV1V2ToV3FormatConverterTest.java | 2 +-
.../local/segment/index/loader/LoaderTest.java | 2 +-
.../SegmentMetadataVirtualColumnProviderTest.java | 2 +-
.../apache/pinot/segment/spi/SegmentMetadata.java | 7 +-
.../spi/index/metadata/SegmentMetadataImpl.java | 8 +-
.../spi/loader/SegmentDirectoryLoaderContext.java | 11 +-
.../pinot/server/api/resources/TablesResource.java | 16 +-
.../server/predownload/PredownloadSegmentInfo.java | 2 +-
.../pinot/server/api/TablesResourceTest.java | 441 ++++++++++-----------
35 files changed, 369 insertions(+), 407 deletions(-)
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/metadata/segment/SegmentZKMetadataUtils.java
b/pinot-common/src/main/java/org/apache/pinot/common/metadata/segment/SegmentZKMetadataUtils.java
index 93a2ad0bbc1..7e4aa5d3b14 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/metadata/segment/SegmentZKMetadataUtils.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/metadata/segment/SegmentZKMetadataUtils.java
@@ -89,7 +89,7 @@ public class SegmentZKMetadataUtils {
segmentZKMetadata.setEndOffset(endOffset);
segmentZKMetadata.setStatus(CommonConstants.Segment.Realtime.Status.DONE);
// for committing segments, we use data CRC to replace but only if the
data CRC is present for the segment
- if (Long.parseLong(segmentMetadata.getDataCrc()) >= 0) {
+ if (segmentMetadata.getDataCrc() >= 0) {
segmentZKMetadata.setUseDataCrc(true);
}
@@ -159,8 +159,8 @@ public class SegmentZKMetadataUtils {
SegmentVersion segmentVersion = segmentMetadata.getVersion();
segmentZKMetadata.setIndexVersion(segmentVersion != null ?
segmentVersion.toString() : null);
segmentZKMetadata.setTotalDocs(segmentMetadata.getTotalDocs());
- segmentZKMetadata.setCrc(Long.parseLong(segmentMetadata.getCrc()));
- segmentZKMetadata.setDataCrc(Long.parseLong(segmentMetadata.getDataCrc()));
+ segmentZKMetadata.setCrc(segmentMetadata.getCrc());
+ segmentZKMetadata.setDataCrc(segmentMetadata.getDataCrc());
segmentZKMetadata.setDownloadUrl(downloadUrl);
segmentZKMetadata.setCrypterName(crypterName);
segmentZKMetadata.setSizeInBytes(segmentSizeInBytes);
diff --git
a/pinot-controller/src/main/java/org/apache/pinot/controller/api/upload/ZKOperator.java
b/pinot-controller/src/main/java/org/apache/pinot/controller/api/upload/ZKOperator.java
index dfe84555d48..4dd9327c045 100644
---
a/pinot-controller/src/main/java/org/apache/pinot/controller/api/upload/ZKOperator.java
+++
b/pinot-controller/src/main/java/org/apache/pinot/controller/api/upload/ZKOperator.java
@@ -300,7 +300,7 @@ public class ZKOperator {
customMapModifierStr != null ? new
SegmentZKMetadataCustomMapModifier(customMapModifierStr) : null;
// Update ZK metadata and refresh the segment if necessary
- long newCrc = Long.parseLong(segmentMetadata.getCrc());
+ long newCrc = segmentMetadata.getCrc();
if (newCrc == existingCrc) {
LOGGER.info(
"New segment crc '{}' is the same as existing segment crc for
segment '{}'. Updating ZK metadata without "
diff --git
a/pinot-controller/src/main/java/org/apache/pinot/controller/recommender/realtime/provisioning/MemoryEstimator.java
b/pinot-controller/src/main/java/org/apache/pinot/controller/recommender/realtime/provisioning/MemoryEstimator.java
index a1d2dc1d97d..dbea3c56c83 100644
---
a/pinot-controller/src/main/java/org/apache/pinot/controller/recommender/realtime/provisioning/MemoryEstimator.java
+++
b/pinot-controller/src/main/java/org/apache/pinot/controller/recommender/realtime/provisioning/MemoryEstimator.java
@@ -404,8 +404,8 @@ public class MemoryEstimator {
segmentZKMetadata.setTimeUnit(segmentMetadata.getTimeUnit());
segmentZKMetadata.setCreationTime(segmentMetadata.getIndexCreationTime());
segmentZKMetadata.setTotalDocs(totalDocs);
- segmentZKMetadata.setCrc(Long.parseLong(segmentMetadata.getCrc()));
- segmentZKMetadata.setDataCrc(Long.parseLong(segmentMetadata.getDataCrc()));
+ segmentZKMetadata.setCrc(segmentMetadata.getCrc());
+ segmentZKMetadata.setDataCrc(segmentMetadata.getDataCrc());
return segmentZKMetadata;
}
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/api/PinotSegmentRestletResourceTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/PinotSegmentRestletResourceTest.java
index dd5008a0fae..5fbd422dad6 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/api/PinotSegmentRestletResourceTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/PinotSegmentRestletResourceTest.java
@@ -245,12 +245,12 @@ public class PinotSegmentRestletResourceTest {
throws Exception {
Map<String, String> crcMap =
adminClient.getSegmentClient().getSegmentToCrcMap(tableName);
if (crcMap == null) {
- crcMap = java.util.Map.of();
+ crcMap = Map.of();
}
for (String segmentName : crcMap.keySet()) {
SegmentMetadata metadata = metadataTable.get(segmentName);
assertNotNull(metadata);
- assertEquals(crcMap.get(segmentName), metadata.getCrc());
+ assertEquals(crcMap.get(segmentName), Long.toString(metadata.getCrc()));
}
assertEquals(crcMap.size(), expectedSize);
}
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/api/upload/ZKOperatorTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/upload/ZKOperatorTest.java
index 339d6b253e6..04b0f9bf90f 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/api/upload/ZKOperatorTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/upload/ZKOperatorTest.java
@@ -158,8 +158,8 @@ public class ZKOperatorTest {
SegmentMetadata segmentMetadata = mock(SegmentMetadata.class);
when(segmentMetadata.getName()).thenReturn(segmentName);
- when(segmentMetadata.getCrc()).thenReturn("12345");
- when(segmentMetadata.getDataCrc()).thenReturn("432");
+ when(segmentMetadata.getCrc()).thenReturn(12345L);
+ when(segmentMetadata.getDataCrc()).thenReturn(432L);
when(segmentMetadata.getIndexCreationTime()).thenReturn(123L);
HttpHeaders httpHeaders = mock(HttpHeaders.class);
@@ -200,8 +200,8 @@ public class ZKOperatorTest {
SegmentMetadata segmentMetadata = mock(SegmentMetadata.class);
when(segmentMetadata.getName()).thenReturn(SEGMENT_NAME);
- when(segmentMetadata.getCrc()).thenReturn("12345");
- when(segmentMetadata.getDataCrc()).thenReturn("432");
+ when(segmentMetadata.getCrc()).thenReturn(12345L);
+ when(segmentMetadata.getDataCrc()).thenReturn(432L);
when(segmentMetadata.getIndexCreationTime()).thenReturn(123L);
HttpHeaders httpHeaders = mock(HttpHeaders.class);
@@ -305,7 +305,7 @@ public class ZKOperatorTest {
assertEquals(segmentZKMetadata.getSizeInBytes(), 10);
// Refresh the segment with a different segment (different CRC)
- when(segmentMetadata.getCrc()).thenReturn("23456");
+ when(segmentMetadata.getCrc()).thenReturn(23456L);
when(segmentMetadata.getIndexCreationTime()).thenReturn(789L);
// Add a tiny sleep to guarantee that refresh time is different from the
previous round
Thread.sleep(10);
@@ -332,8 +332,8 @@ public class ZKOperatorTest {
SegmentMetadata segmentMetadata = mock(SegmentMetadata.class);
when(segmentMetadata.getName()).thenReturn(SEGMENT_NAME);
- when(segmentMetadata.getCrc()).thenReturn("12345");
- when(segmentMetadata.getDataCrc()).thenReturn("432");
+ when(segmentMetadata.getCrc()).thenReturn(12345L);
+ when(segmentMetadata.getDataCrc()).thenReturn(432L);
zkOperator.completeSegmentOperations(REALTIME_TABLE_CONFIG,
segmentMetadata, FileUploadType.SEGMENT, null, null,
"downloadUrl", "downloadUrl", null, 10, true, true,
mock(HttpHeaders.class));
@@ -345,7 +345,7 @@ public class ZKOperatorTest {
// Uploading a segment with LLC segment name but without start/end offset
should fail
when(segmentMetadata.getName()).thenReturn(LLC_SEGMENT_NAME);
- when(segmentMetadata.getCrc()).thenReturn("23456");
+ when(segmentMetadata.getCrc()).thenReturn(23456L);
try {
zkOperator.completeSegmentOperations(REALTIME_TABLE_CONFIG,
segmentMetadata, FileUploadType.SEGMENT, null, null,
"downloadUrl", "downloadUrl", null, 10, true, true,
mock(HttpHeaders.class));
@@ -368,7 +368,7 @@ public class ZKOperatorTest {
assertEquals(segmentZKMetadata.getEndOffset(), "1234");
// Refreshing a segment with LLC segment name but without start/end offset
should success
- when(segmentMetadata.getCrc()).thenReturn("34567");
+ when(segmentMetadata.getCrc()).thenReturn(34567L);
when(segmentMetadata.getStartOffset()).thenReturn(null);
when(segmentMetadata.getEndOffset()).thenReturn(null);
zkOperator.completeSegmentOperations(REALTIME_TABLE_CONFIG,
segmentMetadata, FileUploadType.SEGMENT, null, null,
@@ -381,7 +381,7 @@ public class ZKOperatorTest {
assertEquals(segmentZKMetadata.getEndOffset(), "1234");
// Refreshing a segment with LLC segment name and start/end offset should
override the offsets
- when(segmentMetadata.getCrc()).thenReturn("45678");
+ when(segmentMetadata.getCrc()).thenReturn(45678L);
when(segmentMetadata.getStartOffset()).thenReturn("1234");
when(segmentMetadata.getEndOffset()).thenReturn("2345");
zkOperator.completeSegmentOperations(REALTIME_TABLE_CONFIG,
segmentMetadata, FileUploadType.SEGMENT, null, null,
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/lineage/LineageDeleteInterleavingIntegrationTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/lineage/LineageDeleteInterleavingIntegrationTest.java
index bede9929c2b..f0414f81a71 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/lineage/LineageDeleteInterleavingIntegrationTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/lineage/LineageDeleteInterleavingIntegrationTest.java
@@ -429,7 +429,7 @@ public class LineageDeleteInterleavingIntegrationTest {
}
private void addSegmentWithTime(String tableNameWithType, String
segmentName, long startTimeMs, long endTimeMs) {
- String crc = Long.toString(System.nanoTime());
+ long crc = System.nanoTime();
SegmentMetadata metadata =
SegmentMetadataMockUtils.mockSegmentMetadata(tableNameWithType, segmentName,
100, crc,
startTimeMs, endTimeMs, TimeUnit.MILLISECONDS);
_resourceManager.addNewSegment(tableNameWithType, metadata, "downloadUrl");
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/PinotLLCRealtimeSegmentManagerTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/PinotLLCRealtimeSegmentManagerTest.java
index 039346f46d6..34deec66251 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/PinotLLCRealtimeSegmentManagerTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/PinotLLCRealtimeSegmentManagerTest.java
@@ -143,8 +143,8 @@ public class PinotLLCRealtimeSegmentManagerTest {
private static final long END_TIME_MS = START_TIME_MS +
TimeUnit.HOURS.toMillis(RANDOM.nextInt(24) + 1);
private static final Interval INTERVAL = new Interval(START_TIME_MS,
END_TIME_MS);
// NOTE: CRC is always non-negative
- private static final String CRC = Long.toString(RANDOM.nextLong() &
0xFFFFFFFFL);
- private static final String DATA_CRC = Long.toString(RANDOM.nextLong() &
0xFFFFFFFFL);
+ private static final long CRC = RANDOM.nextLong() & 0xFFFFFFFFL;
+ private static final long DATA_CRC = RANDOM.nextLong() & 0xFFFFFFFFL;
private static final SegmentVersion SEGMENT_VERSION = RANDOM.nextBoolean() ?
SegmentVersion.v1 : SegmentVersion.v3;
@AfterClass
@@ -341,8 +341,8 @@ public class PinotLLCRealtimeSegmentManagerTest {
assertEquals(committedSegmentZKMetadata.getStartOffset(),
PARTITION_OFFSET.toString());
assertEquals(committedSegmentZKMetadata.getEndOffset(), NEXT_OFFSET);
assertEquals(committedSegmentZKMetadata.getCreationTime(),
CURRENT_TIME_MS);
- assertEquals(committedSegmentZKMetadata.getCrc(), Long.parseLong(CRC));
- assertEquals(committedSegmentZKMetadata.getDataCrc(),
Long.parseLong(DATA_CRC));
+ assertEquals(committedSegmentZKMetadata.getCrc(), CRC);
+ assertEquals(committedSegmentZKMetadata.getDataCrc(), DATA_CRC);
assertEquals(committedSegmentZKMetadata.getIndexVersion(),
SEGMENT_VERSION.name());
assertEquals(committedSegmentZKMetadata.getTotalDocs(), NUM_DOCS);
assertEquals(committedSegmentZKMetadata.getSizeInBytes(),
SEGMENT_SIZE_IN_BYTES);
@@ -419,7 +419,7 @@ public class PinotLLCRealtimeSegmentManagerTest {
assertEquals(committedSegmentZKMetadata.getStartOffset(),
committingSegmentStartOffset);
assertEquals(committedSegmentZKMetadata.getEndOffset(),
committingSegmentEndOffset);
assertEquals(committedSegmentZKMetadata.getCreationTime(),
CURRENT_TIME_MS);
- assertEquals(committedSegmentZKMetadata.getCrc(), Long.parseLong(CRC));
+ assertEquals(committedSegmentZKMetadata.getCrc(), CRC);
assertEquals(committedSegmentZKMetadata.getIndexVersion(),
SEGMENT_VERSION.name());
assertEquals(committedSegmentZKMetadata.getTotalDocs(), NUM_DOCS);
assertEquals(committedSegmentZKMetadata.getSizeInBytes(),
SEGMENT_SIZE_IN_BYTES);
@@ -614,7 +614,7 @@ public class PinotLLCRealtimeSegmentManagerTest {
assertEquals(committedSegmentZKMetadata.getStartOffset(),
PARTITION_OFFSET.toString());
assertEquals(committedSegmentZKMetadata.getEndOffset(), NEXT_OFFSET);
assertEquals(committedSegmentZKMetadata.getCreationTime(),
CURRENT_TIME_MS);
- assertEquals(committedSegmentZKMetadata.getCrc(), Long.parseLong(CRC));
+ assertEquals(committedSegmentZKMetadata.getCrc(), CRC);
assertEquals(committedSegmentZKMetadata.getIndexVersion(),
SEGMENT_VERSION.name());
assertEquals(committedSegmentZKMetadata.getTotalDocs(), NUM_DOCS);
assertEquals(committedSegmentZKMetadata.getSizeInBytes(),
SEGMENT_SIZE_IN_BYTES);
@@ -678,7 +678,7 @@ public class PinotLLCRealtimeSegmentManagerTest {
assertEquals(committedSegmentZKMetadata.getStartOffset(),
PARTITION_OFFSET.toString());
assertEquals(committedSegmentZKMetadata.getEndOffset(), NEXT_OFFSET);
assertEquals(committedSegmentZKMetadata.getCreationTime(),
CURRENT_TIME_MS);
- assertEquals(committedSegmentZKMetadata.getCrc(), Long.parseLong(CRC));
+ assertEquals(committedSegmentZKMetadata.getCrc(), CRC);
assertEquals(committedSegmentZKMetadata.getIndexVersion(),
SEGMENT_VERSION.name());
assertEquals(committedSegmentZKMetadata.getTotalDocs(), NUM_DOCS);
assertEquals(committedSegmentZKMetadata.getSizeInBytes(),
SEGMENT_SIZE_IN_BYTES);
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/utils/SegmentMetadataMockUtils.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/utils/SegmentMetadataMockUtils.java
index 532b5dc2b0f..8340c288fc6 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/utils/SegmentMetadataMockUtils.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/utils/SegmentMetadataMockUtils.java
@@ -41,13 +41,13 @@ public class SegmentMetadataMockUtils {
}
public static SegmentMetadata mockSegmentMetadata(String tableName, String
segmentName, int numTotalDocs,
- String crc, long startTime, long endTime, TimeUnit timeUnit) {
+ long crc, long startTime, long endTime, TimeUnit timeUnit) {
SegmentMetadata segmentMetadata = Mockito.mock(SegmentMetadata.class);
Mockito.when(segmentMetadata.getTableName()).thenReturn(tableName);
Mockito.when(segmentMetadata.getName()).thenReturn(segmentName);
Mockito.when(segmentMetadata.getTotalDocs()).thenReturn(numTotalDocs);
Mockito.when(segmentMetadata.getCrc()).thenReturn(crc);
-
Mockito.when(segmentMetadata.getDataCrc()).thenReturn(String.valueOf(Long.parseLong(crc)
+ 100));
+ Mockito.when(segmentMetadata.getDataCrc()).thenReturn(crc + 100);
Mockito.when(segmentMetadata.getStartTime()).thenReturn(startTime);
Mockito.when(segmentMetadata.getEndTime()).thenReturn(endTime);
Mockito.when(segmentMetadata.getTimeInterval()).thenReturn(
@@ -58,29 +58,27 @@ public class SegmentMetadataMockUtils {
}
public static SegmentMetadata mockSegmentMetadata(String tableName, String
segmentName, int numTotalDocs,
- String crc) {
+ long crc) {
return mockSegmentMetadata(tableName, segmentName, numTotalDocs, crc, 1L,
10L, TimeUnit.DAYS);
}
public static SegmentMetadata mockSegmentMetadata(String tableName) {
- String uniqueNumericString = nextUniqueNumericString();
- return mockSegmentMetadata(tableName, tableName + uniqueNumericString,
100, uniqueNumericString);
+ long uniqueId = nextUniqueId();
+ return mockSegmentMetadata(tableName, tableName + uniqueId, 100, uniqueId);
}
public static SegmentMetadata mockSegmentMetadata(String tableName, long
startTime,
long endTime, TimeUnit timeUnit) {
- String uniqueNumericString = nextUniqueNumericString();
- return mockSegmentMetadata(tableName, tableName + uniqueNumericString, 100,
- uniqueNumericString, startTime, endTime, timeUnit);
+ long uniqueId = nextUniqueId();
+ return mockSegmentMetadata(tableName, tableName + uniqueId, 100, uniqueId,
startTime, endTime, timeUnit);
}
public static SegmentMetadata mockSegmentMetadata(String tableName, String
segmentName) {
- String uniqueNumericString = nextUniqueNumericString();
- return mockSegmentMetadata(tableName, segmentName, 100,
uniqueNumericString);
+ return mockSegmentMetadata(tableName, segmentName, 100, nextUniqueId());
}
- private static String nextUniqueNumericString() {
- return Long.toString(UNIQUE_ID_GENERATOR.incrementAndGet());
+ private static long nextUniqueId() {
+ return UNIQUE_ID_GENERATOR.incrementAndGet();
}
public static SegmentZKMetadata mockSegmentZKMetadata(String segmentName,
long numTotalDocs) {
@@ -91,7 +89,7 @@ public class SegmentMetadataMockUtils {
}
public static SegmentMetadata mockSegmentMetadata(String tableName, String
segmentName, int numTotalDocs,
- String crc, long startTime, long endTime, TimeUnit timeUnit, String
partitionColumn, int partitionId,
+ long crc, long startTime, long endTime, TimeUnit timeUnit, String
partitionColumn, int partitionId,
int numPartitions) {
SegmentMetadata segmentMetadata =
mockSegmentMetadata(tableName, segmentName, numTotalDocs, crc,
startTime, endTime, timeUnit);
@@ -117,8 +115,8 @@ public class SegmentMetadataMockUtils {
}
when(segmentMetadata.getTableName()).thenReturn(rawTableName);
when(segmentMetadata.getName()).thenReturn(segmentName);
- when(segmentMetadata.getCrc()).thenReturn("0");
- when(segmentMetadata.getDataCrc()).thenReturn("1");
+ when(segmentMetadata.getCrc()).thenReturn(0L);
+ when(segmentMetadata.getDataCrc()).thenReturn(1L);
TreeMap<String, ColumnMetadata> columnMetadataMap = new TreeMap<>();
columnMetadataMap.put(columnName, columnMetadata);
@@ -131,8 +129,8 @@ public class SegmentMetadataMockUtils {
Mockito.when(segmentMetadata.getTableName()).thenReturn(tableName);
Mockito.when(segmentMetadata.getName()).thenReturn(segmentName);
Mockito.when(segmentMetadata.getTotalDocs()).thenReturn(10);
-
Mockito.when(segmentMetadata.getCrc()).thenReturn(Long.toString(System.nanoTime()));
-
Mockito.when(segmentMetadata.getDataCrc()).thenReturn(Long.toString(System.nanoTime()));
+ Mockito.when(segmentMetadata.getCrc()).thenReturn(System.nanoTime());
+ Mockito.when(segmentMetadata.getDataCrc()).thenReturn(System.nanoTime());
Mockito.when(segmentMetadata.getStartTime()).thenReturn(endTime - 10);
Mockito.when(segmentMetadata.getEndTime()).thenReturn(endTime);
Mockito.when(segmentMetadata.getTimeInterval()).thenReturn(
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/validation/ValidationManagerTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/validation/ValidationManagerTest.java
index 97f4e393873..fec0d0e3b79 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/validation/ValidationManagerTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/validation/ValidationManagerTest.java
@@ -93,7 +93,7 @@ public class ValidationManagerTest {
TEST_INSTANCE.getHelixAdmin().getResourceExternalView(TEST_INSTANCE.getHelixClusterName(),
offlineTableName);
return externalView != null &&
externalView.getPartitionSet().contains(TEST_SEGMENT_NAME);
}, 30_000L, "Failed to find the segment in the ExternalView");
-
Mockito.when(segmentMetadata.getCrc()).thenReturn(Long.toString(System.nanoTime()));
+ Mockito.when(segmentMetadata.getCrc()).thenReturn(System.nanoTime());
TEST_INSTANCE.getHelixResourceManager()
.refreshSegment(offlineTableName, segmentMetadata, segmentZKMetadata,
EXPECTED_VERSION, "downloadUrl");
diff --git
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
index 4da4cf0d6ef..8a3269d6eb1 100644
---
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
+++
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
@@ -1175,7 +1175,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
_logger.info("Reloading existing segment: {} on tier: {}", segmentName,
TierConfigUtils.normalizeTierName(segmentTier));
SegmentDirectory segmentDirectory =
- initSegmentDirectory(segmentName,
String.valueOf(zkMetadata.getCrc()), indexLoadingConfig, zkMetadata);
+ initSegmentDirectory(segmentName, zkMetadata.getCrc(),
indexLoadingConfig, zkMetadata);
// We should first try to reuse existing segment directory
if (canReuseExistingDirectoryForReload(zkMetadata, segmentTier,
segmentDirectory, indexLoadingConfig)) {
_logger.info("Reloading segment: {} using existing segment directory
as no reprocessing needed", segmentName);
@@ -1587,7 +1587,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
// Creates the SegmentDirectory object to access the segment metadata.
// The metadata is null if the segment doesn't exist yet.
SegmentDirectory segmentDirectory =
- tryInitSegmentDirectory(segmentName,
String.valueOf(zkMetadata.getCrc()), indexLoadingConfig, zkMetadata);
+ tryInitSegmentDirectory(segmentName, zkMetadata.getCrc(),
indexLoadingConfig, zkMetadata);
SegmentMetadataImpl segmentMetadata = (segmentDirectory == null) ? null :
segmentDirectory.getSegmentMetadata();
/*
@@ -1630,8 +1630,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
// Close the stale SegmentDirectory object and recreate it with
reprocessed segment.
closeSegmentDirectoryQuietly(segmentDirectory);
ImmutableSegmentLoader.preprocess(indexDir, indexLoadingConfig,
_segmentOperationsThrottlerSet, zkMetadata);
- segmentDirectory = initSegmentDirectory(segmentName,
String.valueOf(zkMetadata.getCrc()),
- indexLoadingConfig, zkMetadata);
+ segmentDirectory = initSegmentDirectory(segmentName,
zkMetadata.getCrc(), indexLoadingConfig, zkMetadata);
}
ImmutableSegment segment = ImmutableSegmentLoader.load(segmentDirectory,
indexLoadingConfig);
_logger.info("Loaded existing segment: {} with CRC: {} on tier: {}",
segmentName, zkMetadata.getCrc(),
@@ -1646,7 +1645,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
}
@Nullable
- protected SegmentDirectory tryInitSegmentDirectory(String segmentName,
String segmentCrc,
+ protected SegmentDirectory tryInitSegmentDirectory(String segmentName, long
segmentCrc,
IndexLoadingConfig indexLoadingConfig, @Nullable SegmentZKMetadata
zkMetadata) {
try {
return initSegmentDirectory(segmentName, segmentCrc, indexLoadingConfig,
zkMetadata);
@@ -1984,7 +1983,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
return new StaleSegment(segmentName, false, null);
}
- protected SegmentDirectory initSegmentDirectory(String segmentName, String
segmentCrc,
+ protected SegmentDirectory initSegmentDirectory(String segmentName, long
segmentCrc,
IndexLoadingConfig indexLoadingConfig, @Nullable SegmentZKMetadata
zkMetadata)
throws Exception {
SegmentDirectoryLoaderContext loaderContext = new
SegmentDirectoryLoaderContext.Builder()
@@ -2009,13 +2008,13 @@ public abstract class BaseTableDataManager implements
TableDataManager {
// CRC check can be performed on both segment CRC and data CRC (if
available) based on the ZK property value of
// useDataCRC.
public static boolean hasSameCRC(SegmentZKMetadata zkMetadata,
SegmentMetadata localMetadata) {
- if (zkMetadata.getCrc() == Long.parseLong(localMetadata.getCrc())) {
+ if (zkMetadata.getCrc() == localMetadata.getCrc()) {
return true;
}
return zkMetadata.isUseDataCrc()
&& zkMetadata.getDataCrc() >= 0
- && Long.parseLong(localMetadata.getDataCrc()) >= 0
- && zkMetadata.getDataCrc() ==
Long.parseLong(localMetadata.getDataCrc());
+ && localMetadata.getDataCrc() >= 0
+ && zkMetadata.getDataCrc() == localMetadata.getDataCrc();
}
protected static void recoverReloadFailureQuietly(String tableNameWithType,
String segmentName, File indexDir) {
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerEnqueueSegmentToReplaceTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerEnqueueSegmentToReplaceTest.java
index 370fd22e0a5..52f0882ec8f 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerEnqueueSegmentToReplaceTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerEnqueueSegmentToReplaceTest.java
@@ -144,7 +144,7 @@ public class
BaseTableDataManagerEnqueueSegmentToReplaceTest {
SegmentMetadata mockMetadata = mock(SegmentMetadata.class);
when(mockSegment.getSegmentMetadata()).thenReturn(mockMetadata);
- when(mockMetadata.getCrc()).thenReturn("12345");
+ when(mockMetadata.getCrc()).thenReturn(12345L);
SegmentZKMetadata mockZkMetadata = mock(SegmentZKMetadata.class);
when(mockZkMetadata.getSegmentName()).thenReturn(SEGMENT_NAME);
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
index 5a5b7de8bdc..c1fe11501e1 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
@@ -162,7 +162,7 @@ public class BaseTableDataManagerTest {
// Mock the case where segment is loaded but its CRC is different from
// the one in zk, thus raw segment is downloaded and loaded.
SegmentMetadata localMetadata = mock(SegmentMetadata.class);
- when(localMetadata.getCrc()).thenReturn("0");
+ when(localMetadata.getCrc()).thenReturn(0L);
BaseTableDataManager tableDataManager = spy(createTableManager());
tableDataManager.registerSegment(SEGMENT_NAME,
createImmutableSegmentDataManager(SEGMENT_NAME, localMetadata));
@@ -185,7 +185,7 @@ public class BaseTableDataManagerTest {
// Mock the case where segment is loaded but its CRC is different from
// the one in zk, thus raw segment is downloaded and loaded.
SegmentMetadata localMetadata = mock(SegmentMetadata.class);
- when(localMetadata.getCrc()).thenReturn("0");
+ when(localMetadata.getCrc()).thenReturn(0L);
// No dataDir for coolTier, thus stay on default tier.
BaseTableDataManager tableDataManager = spy(createTableManager());
@@ -222,7 +222,7 @@ public class BaseTableDataManagerTest {
SegmentZKMetadata zkMetadata = mock(SegmentZKMetadata.class);
when(zkMetadata.getCrc()).thenReturn(crc);
SegmentMetadata localMetadata = mock(SegmentMetadata.class);
- when(localMetadata.getCrc()).thenReturn(Long.toString(crc));
+ when(localMetadata.getCrc()).thenReturn(crc);
BaseTableDataManager tableDataManager = spy(createTableManager());
tableDataManager.registerSegment(SEGMENT_NAME,
createImmutableSegmentDataManager(SEGMENT_NAME, localMetadata));
@@ -253,7 +253,7 @@ public class BaseTableDataManagerTest {
when(zkMetadata.getCrc()).thenReturn(crc);
when(zkMetadata.getTier()).thenReturn(TIER_NAME);
SegmentMetadata localMetadata = mock(SegmentMetadata.class);
- when(localMetadata.getCrc()).thenReturn(Long.toString(crc));
+ when(localMetadata.getCrc()).thenReturn(crc);
// No dataDir for coolTier, thus stay on default tier.
BaseTableDataManager tableDataManager = spy(createTableManager());
@@ -289,7 +289,7 @@ public class BaseTableDataManagerTest {
SegmentZKMetadata zkMetadata = mock(SegmentZKMetadata.class);
when(zkMetadata.getCrc()).thenReturn(crc);
SegmentMetadata localMetadata = mock(SegmentMetadata.class);
- when(localMetadata.getCrc()).thenReturn(Long.toString(crc));
+ when(localMetadata.getCrc()).thenReturn(crc);
// Require to use v3 format.
IndexLoadingConfig indexLoadingConfig = new IndexLoadingConfig();
@@ -319,7 +319,7 @@ public class BaseTableDataManagerTest {
SegmentZKMetadata zkMetadata = mock(SegmentZKMetadata.class);
when(zkMetadata.getCrc()).thenReturn(crc);
SegmentMetadata localMetadata = mock(SegmentMetadata.class);
- when(localMetadata.getCrc()).thenReturn(Long.toString(crc));
+ when(localMetadata.getCrc()).thenReturn(crc);
// Require to add indices.
TableConfig tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME)
@@ -346,7 +346,7 @@ public class BaseTableDataManagerTest {
makeRawSegment(indexDir, new File(TEMP_DIR, SEGMENT_NAME +
TarCompressionUtils.TAR_COMPRESSED_FILE_EXTENSION),
false);
SegmentMetadataImpl segmentMetadata = new SegmentMetadataImpl(indexDir);
- assertEquals(Long.parseLong(segmentMetadata.getCrc()),
zkMetadata.getCrc());
+ assertEquals(segmentMetadata.getCrc(), zkMetadata.getCrc());
// Same CRC but force to download.
BaseTableDataManager tableDataManager = spy(createTableManager());
@@ -366,7 +366,7 @@ public class BaseTableDataManagerTest {
tableDataManager.reloadSegment(SEGMENT_NAME, true, null);
assertTrue(indexDir.exists());
segmentMetadata = new SegmentMetadataImpl(indexDir);
- assertEquals(Long.parseLong(segmentMetadata.getCrc()),
zkMetadata.getCrc());
+ assertEquals(segmentMetadata.getCrc(), zkMetadata.getCrc());
assertEquals(segmentMetadata.getTotalDocs(), 5);
}
@@ -847,7 +847,7 @@ public class BaseTableDataManagerTest {
ImmutableSegmentDataManager segmentDataManager =
createImmutableSegmentDataManager(SEGMENT_NAME, 1024L);
SegmentMetadata segmentMetadata =
segmentDataManager.getSegment().getSegmentMetadata();
- when(segmentMetadata.getDataCrc()).thenReturn("99999");
+ when(segmentMetadata.getDataCrc()).thenReturn(99999L);
BaseTableDataManager tableDataManager = createTableManager();
File dataDir = tableDataManager.getSegmentDataDir(SEGMENT_NAME);
@@ -871,7 +871,7 @@ public class BaseTableDataManagerTest {
ImmutableSegmentDataManager segmentDataManager =
createImmutableSegmentDataManager(SEGMENT_NAME, segmentCrc);
SegmentMetadata segmentMetadata =
segmentDataManager.getSegment().getSegmentMetadata();
- when(segmentMetadata.getDataCrc()).thenReturn("11111");
+ when(segmentMetadata.getDataCrc()).thenReturn(11111L);
BaseTableDataManager tableDataManager = createTableManager();
File dataDir = tableDataManager.getSegmentDataDir(SEGMENT_NAME);
@@ -893,7 +893,7 @@ public class BaseTableDataManagerTest {
ImmutableSegmentDataManager segmentDataManager =
createImmutableSegmentDataManager(SEGMENT_NAME, 1024L);
SegmentMetadata segmentMetadata =
segmentDataManager.getSegment().getSegmentMetadata();
- when(segmentMetadata.getDataCrc()).thenReturn("99999");
+ when(segmentMetadata.getDataCrc()).thenReturn(99999L);
BaseTableDataManager tableDataManager = createTableManager();
File dataDir = tableDataManager.getSegmentDataDir(SEGMENT_NAME);
@@ -1182,7 +1182,7 @@ public class BaseTableDataManagerTest {
protected ImmutableSegmentDataManager
createImmutableSegmentDataManager(String segmentName, long crc) {
SegmentMetadata segmentMetadata = mock(SegmentMetadata.class);
- when(segmentMetadata.getCrc()).thenReturn(Long.toString(crc));
+ when(segmentMetadata.getCrc()).thenReturn(crc);
return createImmutableSegmentDataManager(segmentName, segmentMetadata);
}
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManagerTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManagerTest.java
index 57bbc31018f..40993e9564e 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManagerTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManagerTest.java
@@ -122,7 +122,7 @@ public class DimensionTableDataManagerTest {
_indexDir = new File(tableDataDir, segmentName);
SegmentMetadata segmentMetadata = new SegmentMetadataImpl(_indexDir);
_segmentZKMetadata = new SegmentZKMetadata(segmentName);
- _segmentZKMetadata.setCrc(Long.parseLong(segmentMetadata.getCrc()));
+ _segmentZKMetadata.setCrc(segmentMetadata.getCrc());
}
@AfterMethod(alwaysRun = true)
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
index 8843f289cee..fb75d2ef714 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
@@ -898,7 +898,7 @@ public class RealtimeSegmentDataManagerTest {
driver.init(generatorConfig, recordReader);
driver.build();
}
- return Long.parseLong(new SegmentMetadataImpl(new File(resourceDir,
segmentName)).getCrc());
+ return new SegmentMetadataImpl(new File(resourceDir,
segmentName)).getCrc();
}
@Test
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 2da32aa7ff7..72311626f07 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
@@ -153,7 +153,7 @@ public abstract class BaseSingleSegmentConversionExecutor
extends BaseTaskExecut
boolean reuseExistingSegment = false;
if (copyToDeepStore) {
segmentMetadataTarFile =
createSegmentMetadataTarFile(convertedSegmentDir, tempDataDir, segmentName);
- long convertedSegmentCrc = Long.parseLong(new
SegmentMetadataImpl(convertedSegmentDir).getCrc());
+ long convertedSegmentCrc = new
SegmentMetadataImpl(convertedSegmentDir).getCrc();
reuseExistingSegment =
convertedSegmentCrc == Long.parseLong(originalSegmentCrc) &&
StringUtils.isNotEmpty(downloadURL);
}
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 08edf4cecce..292833042ab 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
@@ -329,7 +329,8 @@ public class MinionTaskUtils {
/// Returns the validDocIds bitmap for the segment, resolved across the
servers hosting it. A server contributes
/// its bitmap only when its CRC matches `expectedCrc`, or - when segment
CRCs differ - when `expectedDataCrc`
- /// matches (see [#crcMatches]). `comparisonModeStr` selects the resolution
mode:
+ /// matches (see [#crcMatches]). A negative `expectedDataCrc` means the data
CRC is unavailable and cannot
+ /// resolve a mismatch. `comparisonModeStr` selects the resolution mode:
/// - `UNSAFE`: first usable bitmap; failing / mismatched / non-READY
servers are skipped.
/// - `EQUAL` (default): the bitmap every server agrees on.
/// - `MOST_VALID_DOCS`: the bitmap with the highest valid-doc count.
@@ -339,7 +340,7 @@ public class MinionTaskUtils {
/// for other fetch failures, CRC mismatches, non-GOOD status, or
`EQUAL`-mode consensus failures.
@Nullable
public static RoaringBitmap getValidDocIdFromServerMatchingCrc(String
tableNameWithType, String segmentName,
- String validDocIdsType, MinionContext minionContext, String expectedCrc,
@Nullable String expectedDataCrc,
+ String validDocIdsType, MinionContext minionContext, long expectedCrc,
long expectedDataCrc,
String comparisonModeStr) {
MinionConstants.ValidDocIdsConsensusMode consensusMode =
parseValidDocIdsConsensusMode(comparisonModeStr);
String clusterName = minionContext.getHelixManager().getClusterName();
@@ -380,8 +381,7 @@ public class MinionTaskUtils {
// offheap upsert is used because we will need to delete & add all
primary keys.
// `BaseSingleSegmentConversionExecutor.executeTask()` already checks
for the crc from the task generator
// against the crc from the current segment zk metadata, so we don't
need to check that here.
- if (!crcMatches(parseCrc(expectedCrc), parseCrc(expectedDataCrc),
parseCrc(crcFromValidDocIdsBitmap),
- serverDataCrc)) {
+ if (!crcMatches(expectedCrc, expectedDataCrc,
parseCrc(crcFromValidDocIdsBitmap), serverDataCrc)) {
if (consensusMode == MinionConstants.ValidDocIdsConsensusMode.UNSAFE) {
LOGGER.warn("CRC mismatch for segment: {} from endpoint {},
skipping", segmentName, endpoint);
continue;
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 2d04d1edaec..8539606a78c 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
@@ -63,11 +63,11 @@ public class UpsertCompactionTaskExecutor extends
BaseSingleSegmentConversionExe
String validDocIdsTypeStr =
MinionTaskUtils.getValidDocIdsType(tableConfig.getUpsertConfig(), configs,
UpsertCompactionTask.VALID_DOC_IDS_TYPE).toString();
SegmentMetadataImpl segmentMetadata = new SegmentMetadataImpl(indexDir);
- String originalSegmentCrcFromTaskGenerator =
configs.get(MinionConstants.ORIGINAL_SEGMENT_CRC_KEY);
- String crcFromDeepStorageSegment = segmentMetadata.getCrc();
+ long originalSegmentCrcFromTaskGenerator =
Long.parseLong(configs.get(MinionConstants.ORIGINAL_SEGMENT_CRC_KEY));
+ long 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)) {
+ if (!ignoreCrcMismatch && originalSegmentCrcFromTaskGenerator !=
crcFromDeepStorageSegment) {
String message = "Crc mismatched between ZK and deepstore copy of
segment: " + segmentName
+ ". Expected crc from ZK: " + originalSegmentCrcFromTaskGenerator +
", crc from deepstore: "
+ crcFromDeepStorageSegment;
diff --git
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutor.java
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutor.java
index b13ea3b5e19..d62a2102ecb 100644
---
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutor.java
+++
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutor.java
@@ -20,9 +20,9 @@ package
org.apache.pinot.plugin.minion.tasks.upsertcompactmerge;
import java.io.File;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
import java.util.Map;
-import java.util.Objects;
import java.util.Set;
import java.util.TreeMap;
import java.util.stream.Collectors;
@@ -102,8 +102,9 @@ public class UpsertCompactMergeTaskExecutor extends
BaseMultipleSegmentsConversi
long maxCreationTimeOfMergingSegments =
getMaxZKCreationTimeFromConfig(configs);
// validate if crc of deepstore copies is same as that in ZK of segments
- List<String> originalSegmentCrcFromTaskGenerator =
-
List.of(configs.get(MinionConstants.ORIGINAL_SEGMENT_CRC_KEY).split(","));
+ List<Long> originalSegmentCrcFromTaskGenerator =
+
Arrays.stream(configs.get(MinionConstants.ORIGINAL_SEGMENT_CRC_KEY).split(",")).map(Long::parseLong)
+ .collect(Collectors.toList());
validateCRCForInputSegments(segmentMetadataList,
originalSegmentCrcFromTaskGenerator);
// Executor-only: read comparison mode string from task config (no auth
resolution or URL hits).
@@ -199,10 +200,10 @@ public class UpsertCompactMergeTaskExecutor extends
BaseMultipleSegmentsConversi
return partitionIDSet.iterator().next();
}
- void validateCRCForInputSegments(List<SegmentMetadataImpl>
segmentMetadataList, List<String> expectedCRCList) {
+ void validateCRCForInputSegments(List<SegmentMetadataImpl>
segmentMetadataList, List<Long> expectedCRCList) {
for (int i = 0; i < segmentMetadataList.size(); i++) {
SegmentMetadataImpl segmentMetadata = segmentMetadataList.get(i);
- if (!Objects.equals(segmentMetadata.getCrc(), expectedCRCList.get(i))) {
+ if (segmentMetadata.getCrc() != expectedCRCList.get(i)) {
String message = String.format("Crc mismatched between ZK and
deepstore copy of segment: %s. Expected crc "
+ "from ZK: %s, crc from deepstore: %s",
segmentMetadata.getName(), expectedCRCList.get(i),
segmentMetadata.getCrc());
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 e583d1af432..e677ebfcb10 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
@@ -107,7 +107,7 @@ public class BaseSingleSegmentConversionExecutorTest {
driver.init(config, new GenericRowRecordReader(rows));
driver.build();
_segmentIndexDir = new File(SEGMENT_DIR, SEGMENT_NAME);
- _segmentCrc = Long.parseLong(new
SegmentMetadataImpl(_segmentIndexDir).getCrc());
+ _segmentCrc = new SegmentMetadataImpl(_segmentIndexDir).getCrc();
Assert.assertTrue(DATA_DIR.mkdirs());
MinionContext.getInstance().setDataDir(DATA_DIR);
@@ -259,7 +259,7 @@ public class BaseSingleSegmentConversionExecutorTest {
Set.of(V1Constants.MetadataKeys.METADATA_FILE_NAME,
V1Constants.SEGMENT_CREATION_META));
SegmentMetadataImpl pushedMetadata = new
SegmentMetadataImpl(untarredMetadataDir);
Assert.assertEquals(pushedMetadata.getName(), SEGMENT_NAME);
- Assert.assertEquals(Long.parseLong(pushedMetadata.getCrc()), _segmentCrc);
+ Assert.assertEquals(pushedMetadata.getCrc(), _segmentCrc);
}
/// A plain-path output dir (local deep store) must reach the controller as
a file URI.
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 31632297f51..155bf2f1a6d 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
@@ -383,8 +383,8 @@ public class MinionTaskUtilsTest {
@Test
public void testExecutorDataCrcFallbackMatch() {
List<Object> responses = List.of(makeResponse("seg1", "2000", "5000",
"server1", makeBitmap(4)));
- RoaringBitmap result =
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
"1000",
- "5000", "UNSAFE", responses, new String[]{"server1"}, this);
+ RoaringBitmap result =
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
1000L,
+ 5000L, "UNSAFE", responses, new String[]{"server1"}, this);
assertNotNull(result);
assertEquals(result.getCardinality(), 4);
}
@@ -392,8 +392,8 @@ public class MinionTaskUtilsTest {
@Test
public void testExecutorDataCrcMismatchSkips() {
List<Object> responses = List.of(makeResponse("seg1", "2000", "9999",
"server1", makeBitmap(4)));
- RoaringBitmap result =
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
"1000",
- "5000", "UNSAFE", responses, new String[]{"server1"}, this);
+ RoaringBitmap result =
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
1000L,
+ 5000L, "UNSAFE", responses, new String[]{"server1"}, this);
assertNull(result);
}
@@ -401,7 +401,7 @@ public class MinionTaskUtilsTest {
public void testFetchFailurePreservesNotFoundException() {
List<Object> responses = List.of(new NotFoundException("HTTP 404 Not
Found"));
NotFoundException e = expectThrows(NotFoundException.class,
- () ->
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
"1000", "EQUAL",
+ () ->
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
1000L, "EQUAL",
responses, new String[]{"server1"}, this));
assertTrue(e.getMessage().contains("seg1"), e.getMessage());
assertTrue(e.getMessage().contains("localhost"), e.getMessage());
@@ -412,7 +412,7 @@ public class MinionTaskUtilsTest {
public void testFetchFailureWrapsOtherExceptions() {
List<Object> responses = List.of(new RuntimeException("connection reset"));
IllegalStateException e = expectThrows(IllegalStateException.class,
- () ->
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
"1000", "EQUAL",
+ () ->
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
1000L, "EQUAL",
responses, new String[]{"server1"}, this));
assertTrue(e.getMessage().contains("seg1"), e.getMessage());
assertEquals(e.getCause().getMessage(), "connection reset");
@@ -422,8 +422,8 @@ public class MinionTaskUtilsTest {
public void testUnsafeModeSkipsFetchFailure() {
List<Object> responses = List.of(
new NotFoundException("HTTP 404 Not Found"),
- makeResponse("seg1", "1000", "server2", makeBitmap(3)));
- RoaringBitmap result =
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
"1000",
+ makeResponse("seg1", 1000L, "server2", makeBitmap(3)));
+ RoaringBitmap result =
getValidDocIdFromServerMatchingCrcWithMockedReader("myTable_REALTIME", "seg1",
1000L,
"UNSAFE", responses, new String[]{"server1", "server2"}, this);
assertNotNull(result);
assertEquals(result.getCardinality(), 3);
@@ -438,9 +438,9 @@ public class MinionTaskUtilsTest {
}
/// Builds a ValidDocIdsBitmapResponse for testing: same segmentCrc and GOOD
status.
- private static ValidDocIdsBitmapResponse makeResponse(String segmentName,
String crc, String instanceId,
+ private static ValidDocIdsBitmapResponse makeResponse(String segmentName,
long crc, String instanceId,
RoaringBitmap bitmap) {
- return new ValidDocIdsBitmapResponse(segmentName, crc, null,
ValidDocIdsType.SNAPSHOT,
+ return new ValidDocIdsBitmapResponse(segmentName, Long.toString(crc),
null, ValidDocIdsType.SNAPSHOT,
RoaringBitmapUtils.serialize(bitmap), instanceId,
ServiceStatus.Status.GOOD);
}
@@ -484,14 +484,14 @@ public class MinionTaskUtilsTest {
/// getValidDocIdsBitmapFromServer returns the next element of
responseOrThrowByCallOrder; if it is an Exception,
/// that exception is thrown (simulating fetch failure).
private static RoaringBitmap
getValidDocIdFromServerMatchingCrcWithMockedReader(String tableName,
- String segmentName, String expectedCrc, String consensusMode,
List<Object> responseOrThrowByCallOrder,
+ String segmentName, long expectedCrc, String consensusMode, List<Object>
responseOrThrowByCallOrder,
String[] servers, MinionTaskUtilsTest testInstance) {
- return getValidDocIdFromServerMatchingCrcWithMockedReader(tableName,
segmentName, expectedCrc, null, consensusMode,
+ return getValidDocIdFromServerMatchingCrcWithMockedReader(tableName,
segmentName, expectedCrc, -1, consensusMode,
responseOrThrowByCallOrder, servers, testInstance);
}
private static RoaringBitmap
getValidDocIdFromServerMatchingCrcWithMockedReader(String tableName,
- String segmentName, String expectedCrc, String expectedDataCrc, String
consensusMode,
+ String segmentName, long expectedCrc, long expectedDataCrc, String
consensusMode,
List<Object> responseOrThrowByCallOrder, String[] servers,
MinionTaskUtilsTest testInstance) {
testInstance.setupMinionContextWithServers(tableName, segmentName,
servers);
// Shared across all mock instances (production creates one reader per
server).
@@ -520,7 +520,7 @@ public class MinionTaskUtilsTest {
public void testSameValidDocsEqualConsensus() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(5)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(5)),
@@ -535,7 +535,7 @@ public class MinionTaskUtilsTest {
public void testSameValidDocsMaxValidDocs() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(5)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(5)),
@@ -550,7 +550,7 @@ public class MinionTaskUtilsTest {
public void testSameValidDocsNone() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(5)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(5)),
@@ -565,7 +565,7 @@ public class MinionTaskUtilsTest {
public void testDifferentValidDocsMaxValidDocsMax() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(5)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(3)),
@@ -580,7 +580,7 @@ public class MinionTaskUtilsTest {
public void testsomeServersNoValidDocsEqualConsensus() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(0)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(0)),
@@ -594,7 +594,7 @@ public class MinionTaskUtilsTest {
public void testsomeServersNoValidDocsMaxValidDocs() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(0)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(0)),
@@ -609,7 +609,7 @@ public class MinionTaskUtilsTest {
public void testSomeServersNoValidDocsNone() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(0)),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(0)),
@@ -626,7 +626,7 @@ public class MinionTaskUtilsTest {
public void testOneServerFailsEqualConsensus() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
makeResponse(segmentName, expectedCrc, "server1", makeBitmap(5)),
new RuntimeException("simulated fetch failure"),
@@ -640,7 +640,7 @@ public class MinionTaskUtilsTest {
public void testOneServerFailsNone() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(
new RuntimeException("simulated fetch failure"),
makeResponse(segmentName, expectedCrc, "server2", makeBitmap(3)),
@@ -655,7 +655,7 @@ public class MinionTaskUtilsTest {
public void testAllServersFailMostValidDocs() {
String tableName = "myTable_REALTIME";
String segmentName = "seg1";
- String expectedCrc = "crc1";
+ long expectedCrc = 1000L;
List<Object> responses = List.of(new RuntimeException("simulated"), new
RuntimeException("simulated"),
new RuntimeException("simulated"));
expectThrows(IllegalStateException.class,
diff --git
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutorTest.java
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutorTest.java
index e64b0d0dbbd..fbb03fb3028 100644
---
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutorTest.java
+++
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskExecutorTest.java
@@ -81,11 +81,11 @@ public class UpsertCompactMergeTaskExecutorTest {
SegmentMetadataImpl segment1 = Mockito.mock(SegmentMetadataImpl.class);
SegmentMetadataImpl segment2 = Mockito.mock(SegmentMetadataImpl.class);
- Mockito.when(segment1.getCrc()).thenReturn("1000");
- Mockito.when(segment2.getCrc()).thenReturn("2000");
+ Mockito.when(segment1.getCrc()).thenReturn(1000L);
+ Mockito.when(segment2.getCrc()).thenReturn(2000L);
List<SegmentMetadataImpl> segmentMetadataList = Arrays.asList(segment1,
segment2);
- List<String> expectedCRCList = Arrays.asList("1000", "2000");
+ List<Long> expectedCRCList = List.of(1000L, 2000L);
_taskExecutor.validateCRCForInputSegments(segmentMetadataList,
expectedCRCList);
}
@@ -95,11 +95,11 @@ public class UpsertCompactMergeTaskExecutorTest {
SegmentMetadataImpl segment1 = Mockito.mock(SegmentMetadataImpl.class);
SegmentMetadataImpl segment2 = Mockito.mock(SegmentMetadataImpl.class);
- Mockito.when(segment1.getCrc()).thenReturn("1000");
- Mockito.when(segment2.getCrc()).thenReturn("3000");
+ Mockito.when(segment1.getCrc()).thenReturn(1000L);
+ Mockito.when(segment2.getCrc()).thenReturn(3000L);
List<SegmentMetadataImpl> segmentMetadataList = Arrays.asList(segment1,
segment2);
- List<String> expectedCRCList = Arrays.asList("1000", "2000");
+ List<Long> expectedCRCList = List.of(1000L, 2000L);
_taskExecutor.validateCRCForInputSegments(segmentMetadataList,
expectedCRCList);
}
@@ -202,19 +202,6 @@ public class UpsertCompactMergeTaskExecutorTest {
_taskExecutor.getCommonPartitionIDForSegments(segmentMetadataList);
}
- /// Tests CRC validation with null CRC values.
- @Test(expectedExceptions = IllegalStateException.class)
- public void testValidateCRCForInputSegmentsWithNullCrc() {
- SegmentMetadataImpl segment1 = Mockito.mock(SegmentMetadataImpl.class);
- Mockito.when(segment1.getCrc()).thenReturn(null);
- Mockito.when(segment1.getName()).thenReturn("segment1");
-
- List<SegmentMetadataImpl> segmentMetadataList = Arrays.asList(segment1);
- List<String> expectedCRCList = Arrays.asList("1000");
-
- _taskExecutor.validateCRCForInputSegments(segmentMetadataList,
expectedCRCList);
- }
-
/// Tests handling of empty segment lists.
@Test(expectedExceptions = NoSuchElementException.class)
public void testGetCommonPartitionIDForEmptySegmentList() {
@@ -228,11 +215,11 @@ public class UpsertCompactMergeTaskExecutorTest {
SegmentMetadataImpl segment1 = Mockito.mock(SegmentMetadataImpl.class);
SegmentMetadataImpl segment2 = Mockito.mock(SegmentMetadataImpl.class);
- Mockito.when(segment1.getCrc()).thenReturn("1000");
- Mockito.when(segment2.getCrc()).thenReturn("2000");
+ Mockito.when(segment1.getCrc()).thenReturn(1000L);
+ Mockito.when(segment2.getCrc()).thenReturn(2000L);
List<SegmentMetadataImpl> segmentMetadataList = Arrays.asList(segment1,
segment2);
- List<String> expectedCRCList = Arrays.asList("1000"); // Only one CRC
+ List<Long> expectedCRCList = List.of(1000L); // Only one CRC
_taskExecutor.validateCRCForInputSegments(segmentMetadataList,
expectedCRCList);
}
@@ -270,19 +257,6 @@ public class UpsertCompactMergeTaskExecutorTest {
Assert.assertEquals(result, 1L);
}
- /// Tests CRC validation with whitespace and empty strings.
- @Test(expectedExceptions = IllegalStateException.class)
- public void testValidateCRCWithEmptyString() {
- SegmentMetadataImpl segment1 = Mockito.mock(SegmentMetadataImpl.class);
- Mockito.when(segment1.getCrc()).thenReturn("");
- Mockito.when(segment1.getName()).thenReturn("segment1");
-
- List<SegmentMetadataImpl> segmentMetadataList = Arrays.asList(segment1);
- List<String> expectedCRCList = Arrays.asList("1000");
-
- _taskExecutor.validateCRCForInputSegments(segmentMetadataList,
expectedCRCList);
- }
-
// Helper methods for testing
/// Creates simple test segments (for backward compatibility with existing
tests).
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/data/manager/SegmentDataManager.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/data/manager/SegmentDataManager.java
index d8c45f4c6e6..6d963a14522 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/data/manager/SegmentDataManager.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/data/manager/SegmentDataManager.java
@@ -89,7 +89,7 @@ public abstract class SegmentDataManager {
return List.of(getSegment());
}
- public String getCrc() {
+ public long getCrc() {
return getSegment().getSegmentMetadata().getCrc();
}
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
index 724bd04ef69..8f3f4e62dc7 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
@@ -325,7 +325,7 @@ public class ImmutableSegmentLoader {
segmentVersionOnDisk, segmentVersionToLoad);
}
- private static void preprocessSegment(File indexDir, String segmentName,
String segmentCrc,
+ private static void preprocessSegment(File indexDir, String segmentName,
long segmentCrc,
IndexLoadingConfig indexLoadingConfig, @Nullable
SegmentOperationsThrottlerSet segmentOperationsThrottlerSet,
SegmentZKMetadata zkMetadata)
throws Exception {
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentCrcVirtualColumnProvider.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentCrcVirtualColumnProvider.java
index d64d0683a61..3b4d0a2bfb6 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentCrcVirtualColumnProvider.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentCrcVirtualColumnProvider.java
@@ -24,26 +24,15 @@ import org.apache.pinot.segment.spi.SegmentMetadata;
/// Virtual column provider for `$crc`, the CRC of the segment.
///
-/// The CRC is exposed as a LONG, matching how it is stored everywhere else in
Pinot ([SegmentMetadata#getCrc()]
-/// renders the same `long` as a String). It reads as NULL for CONSUMING
segments, which have no CRC until they are
-/// committed. Grouping by `$segmentName` and `$crc` is a convenient way to
detect replicas of a segment that have
-/// diverged.
+/// The CRC is exposed as a LONG, matching how it is stored everywhere else in
Pinot. It reads as NULL for CONSUMING
+/// segments, which have no CRC until they are committed. Grouping by
`$segmentName` and `$crc` is a convenient way to
+/// detect replicas of a segment that have diverged.
public class SegmentCrcVirtualColumnProvider extends
BaseSegmentMetadataVirtualColumnProvider {
@Nullable
@Override
protected Object extractValue(SegmentMetadata segmentMetadata) {
- String crc = segmentMetadata.getCrc();
- if (crc == null) {
- return null;
- }
- long crcValue;
- try {
- crcValue = Long.parseLong(crc);
- } catch (NumberFormatException e) {
- // The CRC is always rendered from a long, but never fail a segment load
over an unreadable one
- return null;
- }
- // SegmentMetadataImpl renders an unset CRC as Long.MIN_VALUE
- return crcValue != Long.MIN_VALUE ? crcValue : null;
+ // SegmentMetadataImpl reports an unset CRC as Long.MIN_VALUE
+ long crc = segmentMetadata.getCrc();
+ return crc != Long.MIN_VALUE ? crc : null;
}
}
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/realtime/converter/RealtimeSegmentConverterTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/realtime/converter/RealtimeSegmentConverterTest.java
index ca466549b47..da96d42c904 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/realtime/converter/RealtimeSegmentConverterTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/realtime/converter/RealtimeSegmentConverterTest.java
@@ -572,7 +572,7 @@ public class RealtimeSegmentConverterTest implements
PinotBuffersAfterMethodChec
File indexDir = new File(outputDir, segmentName);
SegmentMetadataImpl segmentMetadata = new SegmentMetadataImpl(indexDir);
- assertEquals(segmentMetadata.getCrc(), params[1]);
+ assertEquals(Long.toString(segmentMetadata.getCrc()), params[1]);
assertEquals(segmentMetadata.getVersion(), SegmentVersion.v3);
assertEquals(segmentMetadata.getTotalDocs(), rows.size());
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/SegmentMetadataImplTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/SegmentMetadataImplTest.java
index 2ba291d635c..2ebc1ce8203 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/SegmentMetadataImplTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/SegmentMetadataImplTest.java
@@ -108,7 +108,7 @@ public class SegmentMetadataImplTest {
JsonNode jsonMeta = metadata.toJson(null);
assertEquals(jsonMeta.get("segmentName").asText(), metadata.getName());
- Assert.assertEquals(jsonMeta.get("crc").asLong(),
Long.valueOf(metadata.getCrc()).longValue());
+ Assert.assertEquals(jsonMeta.get("crc").asLong(), metadata.getCrc());
Assert.assertTrue(jsonMeta.get("creatorName").isNull());
assertEquals(jsonMeta.get("creationTimeMillis").asLong(),
metadata.getIndexCreationTime());
assertEquals(jsonMeta.get("timeColumn").asText(),
metadata.getTimeColumn());
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/converter/SegmentV1V2ToV3FormatConverterTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/converter/SegmentV1V2ToV3FormatConverterTest.java
index 7b44a22ddb8..1aff31f39c9 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/converter/SegmentV1V2ToV3FormatConverterTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/converter/SegmentV1V2ToV3FormatConverterTest.java
@@ -114,7 +114,7 @@ public class SegmentV1V2ToV3FormatConverterTest {
Assert.assertFalse(new File(_segmentDirectory,
V1Constants.MetadataKeys.METADATA_FILE_NAME).exists());
SegmentMetadataImpl metaAfterConversion = new
SegmentMetadataImpl(_segmentDirectory);
Assert.assertNotNull(metaAfterConversion);
-
Assert.assertFalse(metaAfterConversion.getCrc().equalsIgnoreCase(String.valueOf(Long.MIN_VALUE)));
+ Assert.assertNotEquals(metaAfterConversion.getCrc(), Long.MIN_VALUE);
Assert.assertEquals(metaAfterConversion.getCrc(),
beforeConversionMeta.getCrc());
Assert.assertTrue(metaAfterConversion.getIndexCreationTime() !=
Long.MIN_VALUE);
Assert.assertEquals(metaAfterConversion.getIndexCreationTime(),
beforeConversionMeta.getIndexCreationTime());
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/LoaderTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/LoaderTest.java
index ac22a9a3e56..c61d4ad7fc6 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/LoaderTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/LoaderTest.java
@@ -230,7 +230,7 @@ public class LoaderTest {
assertEquals(indexSegment.getDataSource(BuiltInVirtualColumn.TOTALDOCS).getDictionary().get(0),
segmentMetadata.getTotalDocs());
assertEquals(indexSegment.getDataSource(BuiltInVirtualColumn.CRC).getDictionary().get(0),
- Long.parseLong(segmentMetadata.getCrc()));
+ segmentMetadata.getCrc());
assertEquals(indexSegment.getDataSource(BuiltInVirtualColumn.CREATIONTIME).getDictionary().get(0),
segmentMetadata.getIndexCreationTime());
for (String column : List.of(BuiltInVirtualColumn.TOTALDOCS,
BuiltInVirtualColumn.CRC,
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentMetadataVirtualColumnProviderTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentMetadataVirtualColumnProviderTest.java
index a336a2e831f..9c0b6664e47 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentMetadataVirtualColumnProviderTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/virtualcolumn/SegmentMetadataVirtualColumnProviderTest.java
@@ -74,7 +74,7 @@ public class SegmentMetadataVirtualColumnProviderTest {
SegmentMetadata segmentMetadata = mock(SegmentMetadata.class);
when(segmentMetadata.getIndexCreationTime()).thenReturn(CREATION_TIME_MS);
when(segmentMetadata.getTimeInterval()).thenReturn(new
Interval(START_TIME_MS, END_TIME_MS, DateTimeZone.UTC));
- when(segmentMetadata.getCrc()).thenReturn(String.valueOf(CRC));
+ when(segmentMetadata.getCrc()).thenReturn(CRC);
when(segmentMetadata.getTotalDocs()).thenReturn(NUM_DOCS);
return segmentMetadata;
}
diff --git
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/SegmentMetadata.java
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/SegmentMetadata.java
index 46257803407..0c9639767e8 100644
---
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/SegmentMetadata.java
+++
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/SegmentMetadata.java
@@ -58,9 +58,12 @@ public interface SegmentMetadata {
Interval getTimeInterval();
- String getCrc();
+ /// Returns the CRC of the whole segment, or `Long.MIN_VALUE` when the
segment has no CRC recorded.
+ long getCrc();
- String getDataCrc();
+ /// Returns the CRC of the segment data only (excluding the metadata), or
`Long.MIN_VALUE` when the segment has no
+ /// data CRC recorded.
+ long getDataCrc();
SegmentVersion getVersion();
diff --git
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/SegmentMetadataImpl.java
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/SegmentMetadataImpl.java
index 31f72dc6506..32a2a668de7 100644
---
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/SegmentMetadataImpl.java
+++
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/SegmentMetadataImpl.java
@@ -370,13 +370,13 @@ public class SegmentMetadataImpl implements
SegmentMetadata {
}
@Override
- public String getCrc() {
- return String.valueOf(_crc);
+ public long getCrc() {
+ return _crc;
}
@Override
- public String getDataCrc() {
- return String.valueOf(_dataCrc);
+ public long getDataCrc() {
+ return _dataCrc;
}
@Override
diff --git
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/loader/SegmentDirectoryLoaderContext.java
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/loader/SegmentDirectoryLoaderContext.java
index 2f29dfa9a4f..6300edd4e9a 100644
---
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/loader/SegmentDirectoryLoaderContext.java
+++
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/loader/SegmentDirectoryLoaderContext.java
@@ -32,13 +32,13 @@ public class SegmentDirectoryLoaderContext {
private final String _instanceId;
private final String _tableDataDir;
private final String _segmentName;
- private final String _segmentCrc;
+ private final long _segmentCrc;
private final String _segmentTier;
private final Map<String, Map<String, String>> _instanceTierConfigs;
private final Map<String, String> _segmentCustomConfigs;
private SegmentDirectoryLoaderContext(ReadMode readMode, TableConfig
tableConfig, Schema schema, String instanceId,
- String tableDataDir, String segmentName, String segmentCrc, String
segmentTier,
+ String tableDataDir, String segmentName, long segmentCrc, String
segmentTier,
Map<String, Map<String, String>> instanceTierConfigs, Map<String,
String> segmentCustomConfigs) {
_readMode = readMode;
_tableConfig = tableConfig;
@@ -76,7 +76,8 @@ public class SegmentDirectoryLoaderContext {
return _segmentName;
}
- public String getSegmentCrc() {
+ /// Returns the CRC of the segment being loaded, or `Long.MIN_VALUE` when
the caller did not supply one.
+ public long getSegmentCrc() {
return _segmentCrc;
}
@@ -99,7 +100,7 @@ public class SegmentDirectoryLoaderContext {
private String _instanceId;
private String _tableDataDir;
private String _segmentName;
- private String _segmentCrc;
+ private long _segmentCrc = Long.MIN_VALUE;
private String _segmentTier;
private Map<String, Map<String, String>> _instanceTierConfigs;
private Map<String, String> _segmentCustomConfigs;
@@ -134,7 +135,7 @@ public class SegmentDirectoryLoaderContext {
return this;
}
- public Builder setSegmentCrc(String segmentCrc) {
+ public Builder setSegmentCrc(long segmentCrc) {
_segmentCrc = segmentCrc;
return this;
}
diff --git
a/pinot-server/src/main/java/org/apache/pinot/server/api/resources/TablesResource.java
b/pinot-server/src/main/java/org/apache/pinot/server/api/resources/TablesResource.java
index e0c474d8904..b4b4c35afee 100644
---
a/pinot-server/src/main/java/org/apache/pinot/server/api/resources/TablesResource.java
+++
b/pinot-server/src/main/java/org/apache/pinot/server/api/resources/TablesResource.java
@@ -554,7 +554,7 @@ public class TablesResource {
try {
Map<String, String> segmentCrcForTable = new HashMap<>();
for (SegmentDataManager segmentDataManager : segmentDataManagers) {
- segmentCrcForTable.put(segmentDataManager.getSegmentName(),
segmentDataManager.getCrc());
+ segmentCrcForTable.put(segmentDataManager.getSegmentName(),
Long.toString(segmentDataManager.getCrc()));
}
return ResourceUtils.convertToJsonString(segmentCrcForTable);
} catch (Exception e) {
@@ -674,7 +674,7 @@ public class TablesResource {
throw new WebApplicationException(msg, Response.Status.NOT_FOUND);
}
byte[] validDocIdsBytes =
RoaringBitmapUtils.serialize(validDocIdSnapshot);
- return new ValidDocIdsBitmapResponse(segmentName,
indexSegment.getSegmentMetadata().getCrc(),
+ return new ValidDocIdsBitmapResponse(segmentName,
Long.toString(indexSegment.getSegmentMetadata().getCrc()),
toReportableDataCrc(indexSegment.getSegmentMetadata().getDataCrc()),
finalValidDocIdsType, validDocIdsBytes,
_serverInstance.getInstanceDataManager().getInstanceId(), status);
} finally {
@@ -767,7 +767,7 @@ public class TablesResource {
validDocIdsMetadata.put("totalDocs", totalDocs);
validDocIdsMetadata.put("totalValidDocs", totalValidDocs);
validDocIdsMetadata.put("totalInvalidDocs", totalInvalidDocs);
- validDocIdsMetadata.put("segmentCrc",
indexSegment.getSegmentMetadata().getCrc());
+ validDocIdsMetadata.put("segmentCrc",
Long.toString(indexSegment.getSegmentMetadata().getCrc()));
String reportableDataCrc =
toReportableDataCrc(indexSegment.getSegmentMetadata().getDataCrc());
if (reportableDataCrc != null) {
validDocIdsMetadata.put("segmentDataCrc", reportableDataCrc);
@@ -801,8 +801,8 @@ public class TablesResource {
/// The segment's data CRC to report, or null when unavailable (negative).
@Nullable
- private static String toReportableDataCrc(String dataCrc) {
- return dataCrc != null && Long.parseLong(dataCrc) >= 0 ? dataCrc : null;
+ private static String toReportableDataCrc(long dataCrc) {
+ return dataCrc >= 0 ? Long.toString(dataCrc) : null;
}
private Pair<ValidDocIdsType, MutableRoaringBitmap>
getValidDocIds(IndexSegment indexSegment,
@@ -973,8 +973,8 @@ public class TablesResource {
String downloadUrl = uploadSegment(segmentTarFile,
realtimeTableNameWithType, segmentName, timeoutMs);
return new TableLLCSegmentUploadResponse(
segmentName,
-
Long.parseLong(segmentDataManager.getSegment().getSegmentMetadata().getCrc()),
-
Long.parseLong(segmentDataManager.getSegment().getSegmentMetadata().getDataCrc()),
+ segmentDataManager.getSegment().getSegmentMetadata().getCrc(),
+ segmentDataManager.getSegment().getSegmentMetadata().getDataCrc(),
downloadUrl);
} finally {
FileUtils.deleteQuietly(segmentTarFile);
@@ -1202,7 +1202,7 @@ public class TablesResource {
String invalidReason = String.format(
"Segment %s is in ONLINE state, but segmentDataManager is
null", segmentName);
return new TableSegmentValidationInfo(false, invalidReason,
-1);
- } else if
(!segmentDataManager.getCrc().equals(String.valueOf(zkMetadata.getCrc()))) {
+ } else if (segmentDataManager.getCrc() != zkMetadata.getCrc()) {
String invalidReason = String.format(
"Segment %s is in ONLINE state, but has CRC mismatch. "
+ "zk_metadata_crc=%s, segment_data_manager_crc=%s",
diff --git
a/pinot-server/src/main/java/org/apache/pinot/server/predownload/PredownloadSegmentInfo.java
b/pinot-server/src/main/java/org/apache/pinot/server/predownload/PredownloadSegmentInfo.java
index cb4810c72d1..cc0c9be81f2 100644
---
a/pinot-server/src/main/java/org/apache/pinot/server/predownload/PredownloadSegmentInfo.java
+++
b/pinot-server/src/main/java/org/apache/pinot/server/predownload/PredownloadSegmentInfo.java
@@ -127,7 +127,7 @@ public class PredownloadSegmentInfo {
.setInstanceId(indexLoadingConfig.getInstanceId())
.setTableDataDir(indexLoadingConfig.getTableDataDir())
.setSegmentName(_segmentName)
- .setSegmentCrc(String.valueOf(_crc))
+ .setSegmentCrc(_crc)
.setSegmentTier(indexLoadingConfig.getSegmentTier())
.setInstanceTierConfigs(indexLoadingConfig.getInstanceTierConfigs())
.build();
diff --git
a/pinot-server/src/test/java/org/apache/pinot/server/api/TablesResourceTest.java
b/pinot-server/src/test/java/org/apache/pinot/server/api/TablesResourceTest.java
index 742614358f9..3cff7dc2355 100644
---
a/pinot-server/src/test/java/org/apache/pinot/server/api/TablesResourceTest.java
+++
b/pinot-server/src/test/java/org/apache/pinot/server/api/TablesResourceTest.java
@@ -62,6 +62,8 @@ import
org.apache.pinot.segment.spi.index.metadata.SegmentMetadataImpl;
import
org.apache.pinot.segment.spi.index.mutable.ThreadSafeMutableRoaringBitmap;
import org.apache.pinot.segment.spi.store.SegmentDirectoryPaths;
import org.apache.pinot.spi.config.table.FieldConfig;
+import org.apache.pinot.spi.config.table.FieldConfig.CompressionCodec;
+import org.apache.pinot.spi.config.table.FieldConfig.EncodingType;
import org.apache.pinot.spi.config.table.IndexingConfig;
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.spi.config.table.TableType;
@@ -73,11 +75,11 @@ import org.apache.pinot.spi.utils.ReadMode;
import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
import org.apache.pinot.spi.utils.builder.TableNameBuilder;
import org.roaringbitmap.buffer.ImmutableRoaringBitmap;
-import org.testng.Assert;
import org.testng.annotations.Test;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
+import static org.testng.Assert.*;
public class TablesResourceTest extends BaseResourceTest {
@@ -91,12 +93,12 @@ public class TablesResourceTest extends BaseResourceTest {
String responseBody = response.readEntity(String.class);
TablesList tablesList = JsonUtils.stringToObject(responseBody,
TablesList.class);
- Assert.assertNotNull(tablesList);
+ assertNotNull(tablesList);
List<String> tables = tablesList.getTables();
- Assert.assertNotNull(tables);
- Assert.assertEquals(tables.size(), 2);
- Assert.assertEquals(tables.get(0), REALTIME_TABLE_NAME);
- Assert.assertEquals(tables.get(1), OFFLINE_TABLE_NAME);
+ assertNotNull(tables);
+ assertEquals(tables.size(), 2);
+ assertEquals(tables.get(0), REALTIME_TABLE_NAME);
+ assertEquals(tables.get(1), OFFLINE_TABLE_NAME);
String secondTable = "secondTable_REALTIME";
addTable(secondTable);
@@ -104,13 +106,13 @@ public class TablesResourceTest extends BaseResourceTest {
responseBody = response.readEntity(String.class);
tablesList = JsonUtils.stringToObject(responseBody, TablesList.class);
- Assert.assertNotNull(tablesList);
+ assertNotNull(tablesList);
tables = tablesList.getTables();
- Assert.assertNotNull(tables);
- Assert.assertEquals(tables.size(), 3);
- Assert.assertTrue(tables.contains(REALTIME_TABLE_NAME));
- Assert.assertTrue(tables.contains(secondTable));
- Assert.assertTrue(tables.contains(OFFLINE_TABLE_NAME));
+ assertNotNull(tables);
+ assertEquals(tables.size(), 3);
+ assertTrue(tables.contains(REALTIME_TABLE_NAME));
+ assertTrue(tables.contains(secondTable));
+ assertTrue(tables.contains(OFFLINE_TABLE_NAME));
}
@Test
@@ -120,25 +122,25 @@ public class TablesResourceTest extends BaseResourceTest {
IndexSegment defaultSegment = _realtimeIndexSegments.get(0);
TableSegments tableSegments =
_webTarget.path(segmentsPath).request().get(TableSegments.class);
- Assert.assertNotNull(tableSegments);
+ assertNotNull(tableSegments);
List<String> segmentNames = tableSegments.getSegments();
- Assert.assertNotNull(segmentNames);
- Assert.assertEquals(segmentNames.size(), 1);
- Assert.assertEquals(segmentNames.get(0),
_realtimeIndexSegments.get(0).getSegmentName());
+ assertNotNull(segmentNames);
+ assertEquals(segmentNames.size(), 1);
+ assertEquals(segmentNames.get(0),
_realtimeIndexSegments.get(0).getSegmentName());
IndexSegment secondSegment = setUpSegment(REALTIME_TABLE_NAME, null, "0",
_realtimeIndexSegments);
tableSegments =
_webTarget.path(segmentsPath).request().get(TableSegments.class);
- Assert.assertNotNull(tableSegments);
+ assertNotNull(tableSegments);
segmentNames = tableSegments.getSegments();
- Assert.assertNotNull(segmentNames);
- Assert.assertEquals(segmentNames.size(), 2);
- Assert.assertTrue(segmentNames.contains(defaultSegment.getSegmentName()));
- Assert.assertTrue(segmentNames.contains(secondSegment.getSegmentName()));
+ assertNotNull(segmentNames);
+ assertEquals(segmentNames.size(), 2);
+ assertTrue(segmentNames.contains(defaultSegment.getSegmentName()));
+ assertTrue(segmentNames.contains(secondSegment.getSegmentName()));
// No such table
Response response =
_webTarget.path("/tables/noSuchTable/segments").request().get(Response.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertNotNull(response);
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -149,8 +151,8 @@ public class TablesResourceTest extends BaseResourceTest {
JsonNode jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(tableIndexesPath).request().get(String.class));
TableIndexMetadataResponse tableIndexMetadataResponse =
JsonUtils.jsonNodeToObject(jsonResponse,
TableIndexMetadataResponse.class);
- Assert.assertNotNull(tableIndexMetadataResponse);
- Assert.assertEquals(tableIndexMetadataResponse.getTotalOnlineSegments(),
_offlineIndexSegments.size());
+ assertNotNull(tableIndexMetadataResponse);
+ assertEquals(tableIndexMetadataResponse.getTotalOnlineSegments(),
_offlineIndexSegments.size());
Map<String, Map<String, Integer>> columnToIndexCountMap = new HashMap<>();
for (ImmutableSegment segment : _offlineIndexSegments) {
@@ -164,12 +166,12 @@ public class TablesResourceTest extends BaseResourceTest {
});
}
- Assert.assertEquals(tableIndexMetadataResponse.getColumnToIndexesCount(),
columnToIndexCountMap);
+ assertEquals(tableIndexMetadataResponse.getColumnToIndexesCount(),
columnToIndexCountMap);
// No such table
Response response =
_webTarget.path("/tables/noSuchTable/indexes").request().get(Response.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertNotNull(response);
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -182,11 +184,11 @@ public class TablesResourceTest extends BaseResourceTest {
JsonNode jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(tableMetadataPath).request().get(String.class));
TableMetadataInfo metadataInfo =
JsonUtils.jsonNodeToObject(jsonResponse, TableMetadataInfo.class);
- Assert.assertNotNull(metadataInfo);
- Assert.assertEquals(metadataInfo.getTableName(), tableNameWithType);
- Assert.assertEquals(metadataInfo.getColumnLengthMap().size(), 0);
- Assert.assertEquals(metadataInfo.getColumnCardinalityMap().size(), 0);
- Assert.assertEquals(metadataInfo.getColumnIndexSizeMap().size(), 0);
+ assertNotNull(metadataInfo);
+ assertEquals(metadataInfo.getTableName(), tableNameWithType);
+ assertEquals(metadataInfo.getColumnLengthMap().size(), 0);
+ assertEquals(metadataInfo.getColumnCardinalityMap().size(), 0);
+ assertEquals(metadataInfo.getColumnIndexSizeMap().size(), 0);
jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(tableMetadataPath)
.queryParam("columns", "column1")
@@ -194,19 +196,17 @@ public class TablesResourceTest extends BaseResourceTest {
.request()
.get(String.class));
metadataInfo = JsonUtils.jsonNodeToObject(jsonResponse,
TableMetadataInfo.class);
- Assert.assertEquals(metadataInfo.getColumnLengthMap().size(), 2);
- Assert.assertEquals(metadataInfo.getColumnCardinalityMap().size(), 2);
- Assert.assertEquals(metadataInfo.getColumnIndexSizeMap().size(), 2);
- Assert.assertTrue(
-
metadataInfo.getColumnIndexSizeMap().get("column1").containsKey(StandardIndexes.dictionary().getId()));
- Assert.assertTrue(
-
metadataInfo.getColumnIndexSizeMap().get("column2").containsKey(StandardIndexes.forward().getId()));
+ assertEquals(metadataInfo.getColumnLengthMap().size(), 2);
+ assertEquals(metadataInfo.getColumnCardinalityMap().size(), 2);
+ assertEquals(metadataInfo.getColumnIndexSizeMap().size(), 2);
+
assertTrue(metadataInfo.getColumnIndexSizeMap().get("column1").containsKey(StandardIndexes.dictionary().getId()));
+
assertTrue(metadataInfo.getColumnIndexSizeMap().get("column2").containsKey(StandardIndexes.forward().getId()));
}
// No such table
Response response =
_webTarget.path("/tables/noSuchTable/metadata").request().get(Response.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertNotNull(response);
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -225,7 +225,7 @@ public class TablesResourceTest extends BaseResourceTest {
try {
String response = _webTarget.path("/tables/" + REALTIME_TABLE_NAME +
"/metadata").request().get(String.class);
TableMetadataInfo metadataInfo = JsonUtils.stringToObject(response,
TableMetadataInfo.class);
- Assert.assertEquals(metadataInfo.getNumSegments(), 2L);
+ assertEquals(metadataInfo.getNumSegments(), 2L);
} finally {
_tableDataManagerMap.put(REALTIME_TABLE_NAME, original);
}
@@ -241,49 +241,49 @@ public class TablesResourceTest extends BaseResourceTest {
JsonNode jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(segmentMetadataPath).request().get(String.class));
SegmentMetadata segmentMetadata = defaultSegment.getSegmentMetadata();
- Assert.assertEquals(jsonResponse.get("segmentName").asText(),
segmentMetadata.getName());
- Assert.assertEquals(jsonResponse.get("crc").asText(),
segmentMetadata.getCrc());
- Assert.assertEquals(jsonResponse.get("creationTimeMillis").asLong(),
segmentMetadata.getIndexCreationTime());
- Assert.assertTrue(jsonResponse.has("startTimeReadable"));
- Assert.assertTrue(jsonResponse.has("endTimeReadable"));
- Assert.assertTrue(jsonResponse.has("creationTimeReadable"));
- Assert.assertEquals(jsonResponse.get("columns").size(), 0);
- Assert.assertEquals(jsonResponse.get("indexes").size(), 0);
+ assertEquals(jsonResponse.get("segmentName").asText(),
segmentMetadata.getName());
+ assertEquals(jsonResponse.get("crc").asLong(), segmentMetadata.getCrc());
+ assertEquals(jsonResponse.get("creationTimeMillis").asLong(),
segmentMetadata.getIndexCreationTime());
+ assertTrue(jsonResponse.has("startTimeReadable"));
+ assertTrue(jsonResponse.has("endTimeReadable"));
+ assertTrue(jsonResponse.has("creationTimeReadable"));
+ assertEquals(jsonResponse.get("columns").size(), 0);
+ assertEquals(jsonResponse.get("indexes").size(), 0);
jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(segmentMetadataPath)
.queryParam("columns", "column1")
.queryParam("columns", "column2")
.request()
.get(String.class));
- Assert.assertEquals(jsonResponse.get("columns").size(), 2);
- Assert.assertEquals(jsonResponse.get("indexes").size(), 2);
-
Assert.assertNotNull(jsonResponse.get("columns").get(0).get("indexSizeMap"));
-
Assert.assertEquals(jsonResponse.get("columns").get(0).get("indexSizeMap").get("forward_index").asText(),
"400008");
-
Assert.assertEquals(jsonResponse.get("columns").get(0).get("indexSizeMap").get("dictionary").asText(),
"206384");
-
Assert.assertNotNull(jsonResponse.get("columns").get(1).get("indexSizeMap"));
-
Assert.assertEquals(jsonResponse.get("columns").get(1).get("indexSizeMap").get("forward_index").asText(),
"400008");
-
Assert.assertEquals(jsonResponse.get("columns").get(1).get("indexSizeMap").get("dictionary").asText(),
"168976");
-
Assert.assertEquals(jsonResponse.get("indexes").get("column1").get("h3-index").asText(),
"NO");
-
Assert.assertEquals(jsonResponse.get("indexes").get("column1").get("fst-index").asText(),
"NO");
-
Assert.assertEquals(jsonResponse.get("indexes").get("column1").get("text-index").asText(),
"NO");
-
Assert.assertEquals(jsonResponse.get("indexes").get("column2").get("h3-index").asText(),
"NO");
-
Assert.assertEquals(jsonResponse.get("indexes").get("column2").get("fst-index").asText(),
"NO");
-
Assert.assertEquals(jsonResponse.get("indexes").get("column2").get("text-index").asText(),
"NO");
+ assertEquals(jsonResponse.get("columns").size(), 2);
+ assertEquals(jsonResponse.get("indexes").size(), 2);
+ assertNotNull(jsonResponse.get("columns").get(0).get("indexSizeMap"));
+
assertEquals(jsonResponse.get("columns").get(0).get("indexSizeMap").get("forward_index").asText(),
"400008");
+
assertEquals(jsonResponse.get("columns").get(0).get("indexSizeMap").get("dictionary").asText(),
"206384");
+ assertNotNull(jsonResponse.get("columns").get(1).get("indexSizeMap"));
+
assertEquals(jsonResponse.get("columns").get(1).get("indexSizeMap").get("forward_index").asText(),
"400008");
+
assertEquals(jsonResponse.get("columns").get(1).get("indexSizeMap").get("dictionary").asText(),
"168976");
+
assertEquals(jsonResponse.get("indexes").get("column1").get("h3-index").asText(),
"NO");
+
assertEquals(jsonResponse.get("indexes").get("column1").get("fst-index").asText(),
"NO");
+
assertEquals(jsonResponse.get("indexes").get("column1").get("text-index").asText(),
"NO");
+
assertEquals(jsonResponse.get("indexes").get("column2").get("h3-index").asText(),
"NO");
+
assertEquals(jsonResponse.get("indexes").get("column2").get("fst-index").asText(),
"NO");
+
assertEquals(jsonResponse.get("indexes").get("column2").get("text-index").asText(),
"NO");
jsonResponse = JsonUtils.stringToJsonNode(
(_webTarget.path(segmentMetadataPath).queryParam("columns",
"*").request().get(String.class)));
int physicalColumnCount = defaultSegment.getPhysicalColumnNames().size();
- Assert.assertEquals(jsonResponse.get("columns").size(),
physicalColumnCount);
- Assert.assertEquals(jsonResponse.get("indexes").size(),
physicalColumnCount);
+ assertEquals(jsonResponse.get("columns").size(), physicalColumnCount);
+ assertEquals(jsonResponse.get("indexes").size(), physicalColumnCount);
Response response = _webTarget.path("/tables/UNKNOWN_TABLE/segments/" +
defaultSegment.getSegmentName())
.request()
.get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
response =
_webTarget.path("/tables/" + REALTIME_TABLE_NAME +
"/segments/UNKNOWN_SEGMENT").request().get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -298,14 +298,14 @@ public class TablesResourceTest extends BaseResourceTest {
.get(String.class));
JsonNode jsonNode = jsonResponse.get(segmentName);
SegmentMetadata segmentMetadata = defaultSegment.getSegmentMetadata();
- Assert.assertEquals(jsonNode.get("segmentName").asText(),
segmentMetadata.getName());
- Assert.assertEquals(jsonNode.get("crc").asText(),
segmentMetadata.getCrc());
- Assert.assertEquals(jsonNode.get("creationTimeMillis").asLong(),
segmentMetadata.getIndexCreationTime());
- Assert.assertTrue(jsonNode.has("startTimeReadable"));
- Assert.assertTrue(jsonNode.has("endTimeReadable"));
- Assert.assertTrue(jsonNode.has("creationTimeReadable"));
- Assert.assertEquals(jsonNode.get("columns").size(), 0);
- Assert.assertEquals(jsonNode.get("indexes").size(), 0);
+ assertEquals(jsonNode.get("segmentName").asText(),
segmentMetadata.getName());
+ assertEquals(jsonNode.get("crc").asLong(), segmentMetadata.getCrc());
+ assertEquals(jsonNode.get("creationTimeMillis").asLong(),
segmentMetadata.getIndexCreationTime());
+ assertTrue(jsonNode.has("startTimeReadable"));
+ assertTrue(jsonNode.has("endTimeReadable"));
+ assertTrue(jsonNode.has("creationTimeReadable"));
+ assertEquals(jsonNode.get("columns").size(), 0);
+ assertEquals(jsonNode.get("indexes").size(), 0);
jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(segmentMetadataPath)
.queryParam("columns", "column1")
@@ -314,16 +314,16 @@ public class TablesResourceTest extends BaseResourceTest {
.request()
.get(String.class));
jsonNode = jsonResponse.get(segmentName);
- Assert.assertEquals(jsonNode.get("columns").size(), 2);
- Assert.assertEquals(jsonNode.get("indexes").size(), 2);
- Assert.assertNotNull(jsonNode.get("columns").get(0).get("indexSizeMap"));
- Assert.assertNotNull(jsonNode.get("columns").get(1).get("indexSizeMap"));
-
Assert.assertEquals(jsonNode.get("indexes").get("column1").get("h3-index").asText(),
"NO");
-
Assert.assertEquals(jsonNode.get("indexes").get("column1").get("fst-index").asText(),
"NO");
-
Assert.assertEquals(jsonNode.get("indexes").get("column1").get("text-index").asText(),
"NO");
-
Assert.assertEquals(jsonNode.get("indexes").get("column2").get("h3-index").asText(),
"NO");
-
Assert.assertEquals(jsonNode.get("indexes").get("column2").get("fst-index").asText(),
"NO");
-
Assert.assertEquals(jsonNode.get("indexes").get("column2").get("text-index").asText(),
"NO");
+ assertEquals(jsonNode.get("columns").size(), 2);
+ assertEquals(jsonNode.get("indexes").size(), 2);
+ assertNotNull(jsonNode.get("columns").get(0).get("indexSizeMap"));
+ assertNotNull(jsonNode.get("columns").get(1).get("indexSizeMap"));
+
assertEquals(jsonNode.get("indexes").get("column1").get("h3-index").asText(),
"NO");
+
assertEquals(jsonNode.get("indexes").get("column1").get("fst-index").asText(),
"NO");
+
assertEquals(jsonNode.get("indexes").get("column1").get("text-index").asText(),
"NO");
+
assertEquals(jsonNode.get("indexes").get("column2").get("h3-index").asText(),
"NO");
+
assertEquals(jsonNode.get("indexes").get("column2").get("fst-index").asText(),
"NO");
+
assertEquals(jsonNode.get("indexes").get("column2").get("text-index").asText(),
"NO");
jsonResponse =
JsonUtils.stringToJsonNode((_webTarget.path(segmentMetadataPath)
.queryParam("columns", "*")
@@ -332,8 +332,8 @@ public class TablesResourceTest extends BaseResourceTest {
.get(String.class)));
int physicalColumnCount = defaultSegment.getPhysicalColumnNames().size();
jsonNode = jsonResponse.get(segmentName);
- Assert.assertEquals(jsonNode.get("columns").size(), physicalColumnCount);
- Assert.assertEquals(jsonNode.get("indexes").size(), physicalColumnCount);
+ assertEquals(jsonNode.get("columns").size(), physicalColumnCount);
+ assertEquals(jsonNode.get("indexes").size(), physicalColumnCount);
}
@Test
@@ -351,8 +351,8 @@ public class TablesResourceTest extends BaseResourceTest {
// Check that crc info is correct
for (ImmutableSegment immutableSegment : immutableSegments) {
String segmentName = immutableSegment.getSegmentName();
- String crc = immutableSegment.getSegmentMetadata().getCrc();
- Assert.assertEquals(segmentsCrc.get(segmentName).asText(), crc);
+ long crc = immutableSegment.getSegmentMetadata().getCrc();
+ assertEquals(segmentsCrc.get(segmentName).asLong(), crc);
}
}
@@ -366,11 +366,11 @@ public class TablesResourceTest extends BaseResourceTest {
// Verify non-existent table and segment download return NOT_FOUND status.
Response response =
_webTarget.path("/tables/UNKNOWN_REALTIME/segments/segmentname").request().get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
response =
_webTarget.path("/tables/" + REALTIME_TABLE_NAME +
"/segments/UNKNOWN_SEGMENT").request().get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -383,13 +383,13 @@ public class TablesResourceTest extends BaseResourceTest {
// Verify non-existent table and segment download return NOT_FOUND status.
Response response =
_webTarget.path("/segments/UNKNOWN_REALTIME/segmentname/validDocIdsBitmap").request().get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
response =
_webTarget.path(String.format("/segments/%s/%s/validDocIdsBitmap",
REALTIME_TABLE_NAME, "UNKNOWN_SEGMENT"))
.request()
.get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -408,21 +408,21 @@ public class TablesResourceTest extends BaseResourceTest {
.post(Entity.json(tableSegments), String.class);
JsonNode validDocIdsMetadata = JsonUtils.stringToJsonNode(response).get(0);
- Assert.assertEquals(validDocIdsMetadata.get("totalDocs").asInt(), 200000);
- Assert.assertEquals(validDocIdsMetadata.get("totalValidDocs").asInt(), 8);
- Assert.assertEquals(validDocIdsMetadata.get("totalInvalidDocs").asInt(),
199992);
- Assert.assertEquals(validDocIdsMetadata.get("segmentCrc").asText(),
segment.getSegmentMetadata().getCrc());
- Assert.assertEquals(validDocIdsMetadata.get("validDocIdsType").asText(),
"SNAPSHOT");
- Assert.assertEquals(validDocIdsMetadata.get("segmentSizeInBytes").asLong(),
+ assertEquals(validDocIdsMetadata.get("totalDocs").asInt(), 200000);
+ assertEquals(validDocIdsMetadata.get("totalValidDocs").asInt(), 8);
+ assertEquals(validDocIdsMetadata.get("totalInvalidDocs").asInt(), 199992);
+ assertEquals(validDocIdsMetadata.get("segmentCrc").asLong(),
segment.getSegmentMetadata().getCrc());
+ assertEquals(validDocIdsMetadata.get("validDocIdsType").asText(),
"SNAPSHOT");
+ assertEquals(validDocIdsMetadata.get("segmentSizeInBytes").asLong(),
((ImmutableSegmentImpl) segment).getSegmentSizeBytes());
- Assert.assertTrue(validDocIdsMetadata.has("segmentCreationTimeMillis"));
-
Assert.assertTrue(validDocIdsMetadata.get("segmentCreationTimeMillis").asLong()
> 0);
+ assertTrue(validDocIdsMetadata.has("segmentCreationTimeMillis"));
+ assertTrue(validDocIdsMetadata.get("segmentCreationTimeMillis").asLong() >
0);
// Verify server status information
- Assert.assertTrue(validDocIdsMetadata.has("serverStatus"), "Server status
should be included in response");
+ assertTrue(validDocIdsMetadata.has("serverStatus"), "Server status should
be included in response");
String serverStatus = validDocIdsMetadata.get("serverStatus").asText();
- Assert.assertNotNull(serverStatus, "Server status should not be null");
- Assert.assertEquals(serverStatus, "NOT_STARTED", serverStatus);
+ assertNotNull(serverStatus, "Server status should not be null");
+ assertEquals(serverStatus, "NOT_STARTED", serverStatus);
}
@Test
@@ -444,21 +444,21 @@ public class TablesResourceTest extends BaseResourceTest {
.post(Entity.json(tableSegments), String.class);
JsonNode validDocIdsMetadata = JsonUtils.stringToJsonNode(response).get(0);
- Assert.assertEquals(validDocIdsMetadata.get("totalDocs").asInt(), 200000);
- Assert.assertEquals(validDocIdsMetadata.get("totalValidDocs").asInt(), 8);
- Assert.assertEquals(validDocIdsMetadata.get("totalInvalidDocs").asInt(),
199992);
- Assert.assertEquals(validDocIdsMetadata.get("segmentCrc").asText(),
segment.getSegmentMetadata().getCrc());
- Assert.assertEquals(validDocIdsMetadata.get("validDocIdsType").asText(),
"SNAPSHOT_WITH_DELETE");
- Assert.assertEquals(validDocIdsMetadata.get("segmentSizeInBytes").asLong(),
+ assertEquals(validDocIdsMetadata.get("totalDocs").asInt(), 200000);
+ assertEquals(validDocIdsMetadata.get("totalValidDocs").asInt(), 8);
+ assertEquals(validDocIdsMetadata.get("totalInvalidDocs").asInt(), 199992);
+ assertEquals(validDocIdsMetadata.get("segmentCrc").asLong(),
segment.getSegmentMetadata().getCrc());
+ assertEquals(validDocIdsMetadata.get("validDocIdsType").asText(),
"SNAPSHOT_WITH_DELETE");
+ assertEquals(validDocIdsMetadata.get("segmentSizeInBytes").asLong(),
((ImmutableSegmentImpl) segment).getSegmentSizeBytes());
- Assert.assertTrue(validDocIdsMetadata.has("segmentCreationTimeMillis"));
-
Assert.assertTrue(validDocIdsMetadata.get("segmentCreationTimeMillis").asLong()
> 0);
+ assertTrue(validDocIdsMetadata.has("segmentCreationTimeMillis"));
+ assertTrue(validDocIdsMetadata.get("segmentCreationTimeMillis").asLong() >
0);
// Verify server status information
- Assert.assertTrue(validDocIdsMetadata.has("serverStatus"), "Server status
should be included in response");
+ assertTrue(validDocIdsMetadata.has("serverStatus"), "Server status should
be included in response");
String serverStatus = validDocIdsMetadata.get("serverStatus").asText();
- Assert.assertNotNull(serverStatus, "Server status should not be null");
- Assert.assertEquals(serverStatus, "NOT_STARTED", serverStatus);
+ assertNotNull(serverStatus, "Server status should not be null");
+ assertEquals(serverStatus, "NOT_STARTED", serverStatus);
}
// Verify metadata file from segments.
@@ -468,7 +468,7 @@ public class TablesResourceTest extends BaseResourceTest {
// Download the segment and save to a temp local file.
Response response =
_webTarget.path(segmentPath).request().get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.OK.getStatusCode());
+ assertEquals(response.getStatus(), Response.Status.OK.getStatusCode());
File segmentFile = response.readEntity(File.class);
File tempMetadataDir = new File(_tempDir, "segment_metadata");
@@ -484,7 +484,7 @@ public class TablesResourceTest extends BaseResourceTest {
// Load segment metadata
SegmentMetadataImpl metadata = new SegmentMetadataImpl(tempMetadataDir);
- Assert.assertEquals(metadata.getTableName(),
TableNameBuilder.extractRawTableName(tableNameWithType));
+ assertEquals(metadata.getTableName(),
TableNameBuilder.extractRawTableName(tableNameWithType));
FileUtils.forceDelete(tempMetadataDir);
}
@@ -520,12 +520,12 @@ public class TablesResourceTest extends BaseResourceTest {
// Check no type (default should be validDocIdsSnapshot)
ValidDocIdsBitmapResponse response =
_webTarget.path(snapshotPath).request().get(ValidDocIdsBitmapResponse.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getSegmentCrc(),
segment.getSegmentMetadata().getCrc());
- Assert.assertEquals(response.getSegmentName(), segment.getSegmentName());
+ assertNotNull(response);
+ assertEquals(response.getSegmentCrc(),
Long.toString(segment.getSegmentMetadata().getCrc()));
+ assertEquals(response.getSegmentName(), segment.getSegmentName());
byte[] validDocIdsSnapshotBitmap = response.getBitmap();
- Assert.assertNotNull(validDocIdsSnapshotBitmap);
- Assert.assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
+ assertNotNull(validDocIdsSnapshotBitmap);
+ assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
validDocIdsSnapshot.getMutableRoaringBitmap());
// Check snapshot type
@@ -533,12 +533,12 @@ public class TablesResourceTest extends BaseResourceTest {
.queryParam("validDocIdsType", ValidDocIdsType.SNAPSHOT.toString())
.request()
.get(ValidDocIdsBitmapResponse.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getSegmentCrc(),
segment.getSegmentMetadata().getCrc());
- Assert.assertEquals(response.getSegmentName(), segment.getSegmentName());
+ assertNotNull(response);
+ assertEquals(response.getSegmentCrc(),
Long.toString(segment.getSegmentMetadata().getCrc()));
+ assertEquals(response.getSegmentName(), segment.getSegmentName());
validDocIdsSnapshotBitmap = response.getBitmap();
- Assert.assertNotNull(validDocIdsSnapshotBitmap);
- Assert.assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
+ assertNotNull(validDocIdsSnapshotBitmap);
+ assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
validDocIdsSnapshot.getMutableRoaringBitmap());
// Check onHeap type
@@ -546,12 +546,12 @@ public class TablesResourceTest extends BaseResourceTest {
.queryParam("validDocIdsType", ValidDocIdsType.IN_MEMORY.toString())
.request()
.get(ValidDocIdsBitmapResponse.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getSegmentCrc(),
segment.getSegmentMetadata().getCrc());
- Assert.assertEquals(response.getSegmentName(), segment.getSegmentName());
+ assertNotNull(response);
+ assertEquals(response.getSegmentCrc(),
Long.toString(segment.getSegmentMetadata().getCrc()));
+ assertEquals(response.getSegmentName(), segment.getSegmentName());
validDocIdsSnapshotBitmap = response.getBitmap();
- Assert.assertNotNull(validDocIdsSnapshotBitmap);
- Assert.assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
+ assertNotNull(validDocIdsSnapshotBitmap);
+ assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
validDocIds.getMutableRoaringBitmap());
// Check onHeapWithDelete type
@@ -559,12 +559,12 @@ public class TablesResourceTest extends BaseResourceTest {
.queryParam("validDocIdsType",
ValidDocIdsType.IN_MEMORY_WITH_DELETE.toString())
.request()
.get(ValidDocIdsBitmapResponse.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getSegmentCrc(),
segment.getSegmentMetadata().getCrc());
- Assert.assertEquals(response.getSegmentName(), segment.getSegmentName());
+ assertNotNull(response);
+ assertEquals(response.getSegmentCrc(),
Long.toString(segment.getSegmentMetadata().getCrc()));
+ assertEquals(response.getSegmentName(), segment.getSegmentName());
validDocIdsSnapshotBitmap = response.getBitmap();
- Assert.assertNotNull(validDocIdsSnapshotBitmap);
- Assert.assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
+ assertNotNull(validDocIdsSnapshotBitmap);
+ assertEquals(new
ImmutableRoaringBitmap(ByteBuffer.wrap(validDocIdsSnapshotBitmap)).toMutableRoaringBitmap(),
queryableDocIds.getMutableRoaringBitmap());
}
@@ -584,11 +584,11 @@ public class TablesResourceTest extends BaseResourceTest {
.request()
.get(ValidDocIdsBitmapResponse.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getSegmentCrc(),
_realtimeIndexSegments.get(0).getSegmentMetadata().getCrc());
- Assert.assertEquals(response.getSegmentName(), segment.getSegmentName());
- Assert.assertEquals(response.getValidDocIdsType(),
ValidDocIdsType.SNAPSHOT_WITH_DELETE);
- Assert.assertNotNull(response.getBitmap());
+ assertNotNull(response);
+ assertEquals(response.getSegmentCrc(),
Long.toString(_realtimeIndexSegments.get(0).getSegmentMetadata().getCrc()));
+ assertEquals(response.getSegmentName(), segment.getSegmentName());
+ assertEquals(response.getValidDocIdsType(),
ValidDocIdsType.SNAPSHOT_WITH_DELETE);
+ assertNotNull(response.getBitmap());
}
@Test
@@ -602,15 +602,15 @@ public class TablesResourceTest extends BaseResourceTest {
String.format("/segments/%s/%s/upload", REALTIME_TABLE_NAME,
LLC_SEGMENT_NAME_FOR_UPLOAD_SUCCESS))
.request()
.post(null);
- Assert.assertEquals(response.getStatus(),
Response.Status.OK.getStatusCode());
- Assert.assertEquals(response.readEntity(String.class),
SEGMENT_DOWNLOAD_URL);
+ assertEquals(response.getStatus(), Response.Status.OK.getStatusCode());
+ assertEquals(response.readEntity(String.class), SEGMENT_DOWNLOAD_URL);
// Verify bad request: table type is offline
response = _webTarget.path(
String.format("/segments/%s/%s/upload", OFFLINE_TABLE_NAME,
_offlineIndexSegments.get(0).getSegmentName()))
.request()
.post(null);
- Assert.assertEquals(response.getStatus(),
Response.Status.BAD_REQUEST.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.BAD_REQUEST.getStatusCode());
// Verify bad request: segment is not low level consumer segment
response = _webTarget.path(
@@ -618,21 +618,21 @@ public class TablesResourceTest extends BaseResourceTest {
_realtimeIndexSegments.get(0).getSegmentName()))
.request()
.post(null);
- Assert.assertEquals(response.getStatus(),
Response.Status.BAD_REQUEST.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.BAD_REQUEST.getStatusCode());
// Verify non-existent segment uploading fail with NOT_FOUND status.
response = _webTarget.path(
String.format("/segments/%s/%s_dummy/upload", RAW_TABLE_NAME,
LLC_SEGMENT_NAME_FOR_UPLOAD_SUCCESS))
.request()
.post(null);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
// Verify fail to upload segment to segment store with internal server
error.
response =
_webTarget.path(String.format("/segments/%s/%s/upload",
RAW_TABLE_NAME, LLC_SEGMENT_NAME_FOR_UPLOAD_FAILURE))
.request()
.post(null);
- Assert.assertEquals(response.getStatus(),
Response.Status.INTERNAL_SERVER_ERROR.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.INTERNAL_SERVER_ERROR.getStatusCode());
}
@Test
@@ -646,38 +646,38 @@ public class TablesResourceTest extends BaseResourceTest {
JsonUtils.stringToJsonNode(_webTarget.path(segmentMetadataPath).request().get(String.class));
SegmentMetadata segmentMetadata = defaultSegment.getSegmentMetadata();
- Assert.assertEquals(jsonResponse.get("segmentName").asText(),
segmentMetadata.getName());
- Assert.assertEquals(jsonResponse.get("crc").asText(),
segmentMetadata.getCrc());
- Assert.assertEquals(jsonResponse.get("creationTimeMillis").asLong(),
segmentMetadata.getIndexCreationTime());
- Assert.assertTrue(jsonResponse.has("startTimeReadable"));
- Assert.assertTrue(jsonResponse.has("endTimeReadable"));
- Assert.assertTrue(jsonResponse.has("creationTimeReadable"));
- Assert.assertEquals(jsonResponse.get("columns").size(), 0);
- Assert.assertEquals(jsonResponse.get("indexes").size(), 0);
+ assertEquals(jsonResponse.get("segmentName").asText(),
segmentMetadata.getName());
+ assertEquals(jsonResponse.get("crc").asLong(), segmentMetadata.getCrc());
+ assertEquals(jsonResponse.get("creationTimeMillis").asLong(),
segmentMetadata.getIndexCreationTime());
+ assertTrue(jsonResponse.has("startTimeReadable"));
+ assertTrue(jsonResponse.has("endTimeReadable"));
+ assertTrue(jsonResponse.has("creationTimeReadable"));
+ assertEquals(jsonResponse.get("columns").size(), 0);
+ assertEquals(jsonResponse.get("indexes").size(), 0);
jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path(segmentMetadataPath)
.queryParam("columns", "column1")
.queryParam("columns", "column2")
.request()
.get(String.class));
- Assert.assertEquals(jsonResponse.get("columns").size(), 2);
- Assert.assertEquals(jsonResponse.get("indexes").size(), 2);
- Assert.assertEquals(jsonResponse.get("star-tree-index").size(), 0);
+ assertEquals(jsonResponse.get("columns").size(), 2);
+ assertEquals(jsonResponse.get("indexes").size(), 2);
+ assertEquals(jsonResponse.get("star-tree-index").size(), 0);
jsonResponse = JsonUtils.stringToJsonNode(
(_webTarget.path(segmentMetadataPath).queryParam("columns",
"*").request().get(String.class)));
int physicalColumnCount = defaultSegment.getPhysicalColumnNames().size();
- Assert.assertEquals(jsonResponse.get("columns").size(),
physicalColumnCount);
- Assert.assertEquals(jsonResponse.get("indexes").size(),
physicalColumnCount);
+ assertEquals(jsonResponse.get("columns").size(), physicalColumnCount);
+ assertEquals(jsonResponse.get("indexes").size(), physicalColumnCount);
Response response = _webTarget.path("/tables/UNKNOWN_TABLE/segments/" +
defaultSegment.getSegmentName())
.request()
.get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
response =
_webTarget.path("/tables/" + REALTIME_TABLE_NAME +
"/segments/UNKNOWN_SEGMENT").request().get(Response.class);
- Assert.assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
+ assertEquals(response.getStatus(),
Response.Status.NOT_FOUND.getStatusCode());
}
@Test
@@ -688,27 +688,29 @@ public class TablesResourceTest extends BaseResourceTest {
addTable(tableName);
TableDataManager tableDataManager = _tableDataManagerMap.get(tableName);
ImmutableSegment trackedSegment = setUpSegment(tableName, null, "tracked",
segments, true);
-
Assert.assertTrue(trackedSegment.getSegmentMetadata().getColumnMetadataMap().values().stream()
+ assertTrue(trackedSegment.getSegmentMetadata()
+ .getColumnMetadataMap()
+ .values()
+ .stream()
.anyMatch(column ->
column.getRawForwardIndexUncompressedValueSizeInBytes() >= 0
|| column.getDictionaryEncodedUncompressedValueSizeInBytes() >=
0));
try {
TableSegments request = new
TableSegments(List.of(trackedSegment.getSegmentName()));
- JsonNode jsonResponse = JsonUtils.stringToJsonNode(_webTarget
- .path("/tables/" + tableName + "/compression-stats")
+ JsonNode jsonResponse =
JsonUtils.stringToJsonNode(_webTarget.path("/tables/" + tableName +
"/compression-stats")
.queryParam("includeColumnCompressionStats", "true")
.request()
.post(Entity.json(request), String.class));
ServerCompressionStatsResponse response =
JsonUtils.jsonNodeToObject(jsonResponse,
ServerCompressionStatsResponse.class);
- Assert.assertNotNull(response);
- Assert.assertEquals(response.getSegmentCompressionStats().size(), 1);
+ assertNotNull(response);
+ assertEquals(response.getSegmentCompressionStats().size(), 1);
SegmentCompressionStatsContribution contribution =
response.getSegmentCompressionStats().get(0);
- Assert.assertFalse(contribution.isComplete());
- Assert.assertEquals(contribution.getUncompressedValueSizeInBytes(), -1);
-
Assert.assertEquals(contribution.getForwardIndexAndDictionaryStorageSizeInBytes(),
-1);
- Assert.assertNull(contribution.getColumnCompressionStats());
+ assertFalse(contribution.isComplete());
+ assertEquals(contribution.getUncompressedValueSizeInBytes(), -1);
+
assertEquals(contribution.getForwardIndexAndDictionaryStorageSizeInBytes(), -1);
+ assertNull(contribution.getColumnCompressionStats());
} finally {
tableDataManager.offloadSegment(trackedSegment.getSegmentName());
tableDataManager.shutDown();
@@ -719,15 +721,15 @@ public class TablesResourceTest extends BaseResourceTest {
@Test
public void testGetCompressionStatsWithMissingSegmentList()
throws Exception {
- JsonNode jsonResponse = JsonUtils.stringToJsonNode(_webTarget
- .path("/tables/" + OFFLINE_TABLE_NAME + "/compression-stats")
- .request()
- .post(Entity.json("{}"), String.class));
+ JsonNode jsonResponse = JsonUtils.stringToJsonNode(
+ _webTarget.path("/tables/" + OFFLINE_TABLE_NAME + "/compression-stats")
+ .request()
+ .post(Entity.json("{}"), String.class));
ServerCompressionStatsResponse response =
JsonUtils.jsonNodeToObject(jsonResponse,
ServerCompressionStatsResponse.class);
- Assert.assertNotNull(response);
- Assert.assertTrue(response.getSegmentCompressionStats().isEmpty());
+ assertNotNull(response);
+ assertTrue(response.getSegmentCompressionStats().isEmpty());
}
@Test
@@ -759,19 +761,15 @@ public class TablesResourceTest extends BaseResourceTest {
SegmentIndexCreationDriverImpl dictDriver = new
SegmentIndexCreationDriverImpl();
dictDriver.init(dictConfig, new GenericRowRecordReader(rows));
dictDriver.build();
- ImmutableSegment dictSegment = ImmutableSegmentLoader.load(
- new File(tableDataDir, dictDriver.getSegmentName()),
- ReadMode.mmap);
+ ImmutableSegment dictSegment =
+ ImmutableSegmentLoader.load(new File(tableDataDir,
dictDriver.getSegmentName()), ReadMode.mmap);
mixedSegments.add(dictSegment);
// Segment 2: raw-encoded for column1 and column2
TableConfig rawTableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(mixedTableName)
.setNoDictionaryColumns(List.of("column1", "column2"))
- .setFieldConfigList(List.of(
- new FieldConfig("column1", FieldConfig.EncodingType.RAW, List.of(),
- FieldConfig.CompressionCodec.LZ4, null),
- new FieldConfig("column2", FieldConfig.EncodingType.RAW, List.of(),
- FieldConfig.CompressionCodec.LZ4, null)))
+ .setFieldConfigList(List.of(new FieldConfig("column1",
EncodingType.RAW, List.of(), CompressionCodec.LZ4, null),
+ new FieldConfig("column2", EncodingType.RAW, List.of(),
CompressionCodec.LZ4, null)))
.build();
rawTableConfig.getIndexingConfig().setCompressionStatsEnabled(true);
SegmentGeneratorConfig rawConfig = new
SegmentGeneratorConfig(rawTableConfig, schema);
@@ -780,17 +778,16 @@ public class TablesResourceTest extends BaseResourceTest {
SegmentIndexCreationDriverImpl rawDriver = new
SegmentIndexCreationDriverImpl();
rawDriver.init(rawConfig, new GenericRowRecordReader(rows));
rawDriver.build();
- ImmutableSegment rawSegment = ImmutableSegmentLoader.load(
- new File(tableDataDir, rawDriver.getSegmentName()),
- ReadMode.mmap);
+ ImmutableSegment rawSegment =
+ ImmutableSegmentLoader.load(new File(tableDataDir,
rawDriver.getSegmentName()), ReadMode.mmap);
for (String column : List.of("column1", "column2")) {
-
Assert.assertFalse(rawSegment.getSegmentMetadata().getColumnMetadataFor(column).hasDictionary());
- Assert.assertEquals(
+
assertFalse(rawSegment.getSegmentMetadata().getColumnMetadataFor(column).hasDictionary());
+ assertEquals(
rawSegment.getSegmentMetadata().getColumnMetadataFor(column).getRawForwardIndexChunkCompressionType(),
ChunkCompressionType.LZ4);
- Assert.assertTrue(
- rawSegment.getSegmentMetadata().getColumnMetadataFor(column)
- .getRawForwardIndexUncompressedValueSizeInBytes() > 0);
+ assertTrue(
+
rawSegment.getSegmentMetadata().getColumnMetadataFor(column).getRawForwardIndexUncompressedValueSizeInBytes()
+ > 0);
}
mixedSegments.add(rawSegment);
@@ -806,48 +803,48 @@ public class TablesResourceTest extends BaseResourceTest {
}
try {
- JsonNode jsonResponse = JsonUtils.stringToJsonNode(_webTarget
- .path("/tables/" + mixedTableName + "/compression-stats")
- .queryParam("columns", "column1")
- .queryParam("columns", "column2")
- .queryParam("includeColumnCompressionStats", "true")
- .request()
- .post(Entity.json(new
TableSegments(List.of(dictSegment.getSegmentName(),
rawSegment.getSegmentName()))),
- String.class));
+ JsonNode jsonResponse = JsonUtils.stringToJsonNode(
+ _webTarget.path("/tables/" + mixedTableName + "/compression-stats")
+ .queryParam("columns", "column1")
+ .queryParam("columns", "column2")
+ .queryParam("includeColumnCompressionStats", "true")
+ .request()
+ .post(Entity.json(new
TableSegments(List.of(dictSegment.getSegmentName(),
rawSegment.getSegmentName()))),
+ String.class));
ServerCompressionStatsResponse compressionResponse =
JsonUtils.jsonNodeToObject(jsonResponse,
ServerCompressionStatsResponse.class);
- Assert.assertNotNull(compressionResponse);
+ assertNotNull(compressionResponse);
for (String column : List.of("column1", "column2")) {
boolean sawDictionary = false;
boolean sawRaw = false;
for (SegmentCompressionStatsContribution segmentStats :
compressionResponse.getSegmentCompressionStats()) {
Map<String, ColumnCompressionStatsContribution> columnStats =
segmentStats.getColumnCompressionStats();
- Assert.assertNotNull(columnStats);
- for (ColumnCompressionStatsContribution.EncodingContribution encoding
- : columnStats.get(column).getEncodingBreakdown()) {
- sawDictionary |= encoding.getEncoding() ==
FieldConfig.EncodingType.DICTIONARY
- && encoding.getChunkCompressionType() == null;
- sawRaw |= encoding.getEncoding() == FieldConfig.EncodingType.RAW
+ assertNotNull(columnStats);
+ for (ColumnCompressionStatsContribution.EncodingContribution
encoding : columnStats.get(column)
+ .getEncodingBreakdown()) {
+ sawDictionary |=
+ encoding.getEncoding() == EncodingType.DICTIONARY &&
encoding.getChunkCompressionType() == null;
+ sawRaw |= encoding.getEncoding() == EncodingType.RAW
&& encoding.getChunkCompressionType() ==
ChunkCompressionType.LZ4;
}
}
- Assert.assertTrue(sawDictionary);
- Assert.assertTrue(sawRaw);
+ assertTrue(sawDictionary);
+ assertTrue(sawRaw);
}
- JsonNode filteredResponse = JsonUtils.stringToJsonNode(_webTarget
- .path("/tables/" + mixedTableName + "/compression-stats")
- .queryParam("columns", "column1")
- .queryParam("includeColumnCompressionStats", "true")
- .request()
- .post(Entity.json(new
TableSegments(List.of(dictSegment.getSegmentName(),
rawSegment.getSegmentName()))),
- String.class));
+ JsonNode filteredResponse = JsonUtils.stringToJsonNode(
+ _webTarget.path("/tables/" + mixedTableName + "/compression-stats")
+ .queryParam("columns", "column1")
+ .queryParam("includeColumnCompressionStats", "true")
+ .request()
+ .post(Entity.json(new
TableSegments(List.of(dictSegment.getSegmentName(),
rawSegment.getSegmentName()))),
+ String.class));
ServerCompressionStatsResponse filteredInfo =
JsonUtils.jsonNodeToObject(filteredResponse,
ServerCompressionStatsResponse.class);
for (SegmentCompressionStatsContribution segmentStats :
filteredInfo.getSegmentCompressionStats()) {
- Assert.assertNotNull(segmentStats.getColumnCompressionStats());
- Assert.assertEquals(segmentStats.getColumnCompressionStats().keySet(),
Set.of("column1"));
+ assertNotNull(segmentStats.getColumnCompressionStats());
+ assertEquals(segmentStats.getColumnCompressionStats().keySet(),
Set.of("column1"));
}
} finally {
for (ImmutableSegment seg : mixedSegments) {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]