This is an automated email from the ASF dual-hosted git repository.
jackie 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 561e8bb286 Add server status in the validDocIds info API - Upsert
Tables (#16165)
561e8bb286 is described below
commit 561e8bb28686a816c728cb82d32d281d6fcd42dd
Author: Chaitanya Deepthi <[email protected]>
AuthorDate: Tue Jun 24 15:30:14 2025 -0700
Add server status in the validDocIds info API - Upsert Tables (#16165)
---
.../resources/ValidDocIdsBitmapResponse.java | 16 +++++-
.../restlet/resources/ValidDocIdsMetadataInfo.java | 16 +++++-
.../helix/core/PinotHelixResourceManager.java | 38 ++++++++++++++
.../pinot/plugin/minion/tasks/MinionTaskUtils.java | 13 +++++
.../UpsertCompactionTaskGenerator.java | 25 ++++++---
.../UpsertCompactMergeTaskGenerator.java | 14 ++++-
.../UpsertCompactionTaskGeneratorTest.java | 14 +++--
.../pinot/server/api/resources/TablesResource.java | 59 ++++++++++++----------
.../apache/pinot/server/api/BaseResourceTest.java | 28 +++++++---
.../pinot/server/api/TablesResourceTest.java | 6 +++
10 files changed, 180 insertions(+), 49 deletions(-)
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsBitmapResponse.java
b/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsBitmapResponse.java
index 9b54774c66..2cb8d66131 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsBitmapResponse.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsBitmapResponse.java
@@ -19,6 +19,7 @@
package org.apache.pinot.common.restlet.resources;
import com.fasterxml.jackson.annotation.JsonProperty;
+import org.apache.pinot.common.utils.ServiceStatus;
public class ValidDocIdsBitmapResponse {
@@ -26,14 +27,19 @@ public class ValidDocIdsBitmapResponse {
private final String _segmentCrc;
private final ValidDocIdsType _validDocIdsType;
private final byte[] _bitmap;
+ private final String _instanceId;
+ private final ServiceStatus.Status _serverStatus;
public ValidDocIdsBitmapResponse(@JsonProperty("segmentName") String
segmentName,
@JsonProperty("segmentCrc") String crc, @JsonProperty("validDocIdsType")
ValidDocIdsType validDocIdsType,
- @JsonProperty("bitmap") byte[] bitmap) {
+ @JsonProperty("bitmap") byte[] bitmap, @JsonProperty("instanceId")
String instanceId,
+ @JsonProperty("serverStatus") ServiceStatus.Status serverStatus) {
_segmentName = segmentName;
_segmentCrc = crc;
_validDocIdsType = validDocIdsType;
_bitmap = bitmap;
+ _instanceId = instanceId;
+ _serverStatus = serverStatus;
}
public String getSegmentName() {
@@ -51,4 +57,12 @@ public class ValidDocIdsBitmapResponse {
public byte[] getBitmap() {
return _bitmap;
}
+
+ public String getInstanceId() {
+ return _instanceId;
+ }
+
+ public ServiceStatus.Status getServerStatus() {
+ return _serverStatus;
+ }
}
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsMetadataInfo.java
b/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsMetadataInfo.java
index 1038d45f55..6fd42c86f6 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsMetadataInfo.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/restlet/resources/ValidDocIdsMetadataInfo.java
@@ -20,6 +20,7 @@ package org.apache.pinot.common.restlet.resources;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonProperty;
+import org.apache.pinot.common.utils.ServiceStatus;
@JsonIgnoreProperties(ignoreUnknown = true)
@@ -32,13 +33,16 @@ public class ValidDocIdsMetadataInfo {
private final ValidDocIdsType _validDocIdsType;
private final long _segmentSizeInBytes;
private final long _segmentCreationTimeMillis;
+ private final String _instanceId;
+ private final ServiceStatus.Status _serverStatus;
public ValidDocIdsMetadataInfo(@JsonProperty("segmentName") String
segmentName,
@JsonProperty("totalValidDocs") long totalValidDocs,
@JsonProperty("totalInvalidDocs") long totalInvalidDocs,
@JsonProperty("totalDocs") long totalDocs, @JsonProperty("segmentCrc")
String segmentCrc,
@JsonProperty("validDocIdsType") ValidDocIdsType validDocIdsType,
@JsonProperty("segmentSizeInBytes") long segmentSizeInBytes,
- @JsonProperty("segmentCreationTimeMillis") long
segmentCreationTimeMillis) {
+ @JsonProperty("segmentCreationTimeMillis") long
segmentCreationTimeMillis,
+ @JsonProperty("instanceId") String instanceId,
@JsonProperty("serverStatus") ServiceStatus.Status serverStatus) {
_segmentName = segmentName;
_totalValidDocs = totalValidDocs;
_totalInvalidDocs = totalInvalidDocs;
@@ -47,6 +51,8 @@ public class ValidDocIdsMetadataInfo {
_validDocIdsType = validDocIdsType;
_segmentSizeInBytes = segmentSizeInBytes;
_segmentCreationTimeMillis = segmentCreationTimeMillis;
+ _instanceId = instanceId;
+ _serverStatus = serverStatus;
}
public String getSegmentName() {
@@ -80,4 +86,12 @@ public class ValidDocIdsMetadataInfo {
public long getSegmentCreationTimeMillis() {
return _segmentCreationTimeMillis;
}
+
+ public String getInstanceId() {
+ return _instanceId;
+ }
+
+ public ServiceStatus.Status getServerStatus() {
+ return _serverStatus;
+ }
}
diff --git
a/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
b/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
index 8b474d2e98..3b0a683970 100644
---
a/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
+++
b/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
@@ -3235,6 +3235,44 @@ public class PinotHelixResourceManager {
return serverToSegmentsMap;
}
+ /**
+ * Get the servers to segments map for which servers are ONLINE in external
view for those segments in IDEAL STATE
+ */
+ public Map<String, List<String>> getServerToOnlineSegmentsMapFromEV(String
tableNameWithType,
+ boolean includeReplacedSegments) {
+ Map<String, List<String>> serverToSegmentsMap = new TreeMap<>();
+ IdealState idealState =
_helixAdmin.getResourceIdealState(_helixClusterName, tableNameWithType);
+ ExternalView externalView =
_helixAdmin.getResourceExternalView(_helixClusterName, tableNameWithType);
+ if (idealState == null) {
+ throw new IllegalStateException("Ideal State does not exist for table: "
+ tableNameWithType);
+ }
+ if (externalView == null) {
+ throw new IllegalStateException("External View state does not exist for
table: " + tableNameWithType);
+ }
+
+ Map<String, Map<String, String>> idealStateMap =
idealState.getRecord().getMapFields();
+ Set<String> segments = idealStateMap.keySet();
+ if (!includeReplacedSegments) {
+ SegmentLineage segmentLineage =
+ SegmentLineageAccessHelper.getSegmentLineage(getPropertyStore(),
tableNameWithType);
+ SegmentLineageUtils.filterSegmentsBasedOnLineageInPlace(segments,
segmentLineage);
+ }
+
+ for (Map.Entry<String, Map<String, String>> entry :
idealStateMap.entrySet()) {
+ String segmentName = entry.getKey();
+ Map<String, String> externalViewStateMap =
externalView.getStateMap(segmentName);
+ if (externalViewStateMap != null) {
+ for (Map.Entry<String, String> instanceStateEntry :
externalViewStateMap.entrySet()) {
+ String server = instanceStateEntry.getKey();
+ if (instanceStateEntry.getValue().equals(SegmentStateModel.ONLINE)) {
+ serverToSegmentsMap.computeIfAbsent(server, key -> new
ArrayList<>()).add(segmentName);
+ }
+ }
+ }
+ }
+ return serverToSegmentsMap;
+ }
+
/**
* Returns a map from server instance to count of segments it serves for the
given table. Ignore OFFLINE segments from
* the ideal state because they are not supposed to be served.
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 e827e25e97..a8be280354 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
@@ -33,6 +33,7 @@ import org.apache.helix.model.ExternalView;
import org.apache.helix.model.InstanceConfig;
import org.apache.pinot.common.restlet.resources.ValidDocIdsBitmapResponse;
import org.apache.pinot.common.utils.RoaringBitmapUtils;
+import org.apache.pinot.common.utils.ServiceStatus;
import org.apache.pinot.common.utils.config.InstanceUtils;
import org.apache.pinot.controller.helix.core.minion.ClusterInfoAccessor;
import org.apache.pinot.controller.util.ServerSegmentMetadataReader;
@@ -236,6 +237,18 @@ public class MinionTaskUtils {
LOGGER.warn(message);
continue;
}
+
+ // skipping servers which are not in READY state. The bitmaps would be
inconsistent when
+ // server is NOT READY as UPDATING segments might be updating the ONLINE
segments
+ if (validDocIdsBitmapResponse.getServerStatus() != null &&
!validDocIdsBitmapResponse.getServerStatus()
+ .equals(ServiceStatus.Status.GOOD)) {
+ String message = "Server " + validDocIdsBitmapResponse.getInstanceId()
+ " is in "
+ + validDocIdsBitmapResponse.getServerStatus() + " state, skipping
it for execution for segment: "
+ + validDocIdsBitmapResponse.getSegmentName() + ". Will try other
servers.";
+ LOGGER.warn(message);
+ continue;
+ }
+
validDocIds =
RoaringBitmapUtils.deserialize(validDocIdsBitmapResponse.getBitmap());
break;
}
diff --git
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java
index 6746225052..ecedff04c0 100644
---
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java
+++
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGenerator.java
@@ -34,6 +34,7 @@ import
org.apache.pinot.common.exception.InvalidConfigException;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo;
import org.apache.pinot.common.restlet.resources.ValidDocIdsType;
+import org.apache.pinot.common.utils.ServiceStatus;
import org.apache.pinot.controller.helix.core.PinotHelixResourceManager;
import
org.apache.pinot.controller.helix.core.minion.generator.BaseTaskGenerator;
import
org.apache.pinot.controller.helix.core.minion.generator.TaskGeneratorUtils;
@@ -124,7 +125,8 @@ public class UpsertCompactionTaskGenerator extends
BaseTaskGenerator {
// get server to segment mappings
PinotHelixResourceManager pinotHelixResourceManager =
_clusterInfoAccessor.getPinotHelixResourceManager();
- Map<String, List<String>> serverToSegments =
pinotHelixResourceManager.getServerToSegmentsMap(tableNameWithType);
+ Map<String, List<String>> serverToSegments =
+
pinotHelixResourceManager.getServerToOnlineSegmentsMapFromEV(tableNameWithType,
true);
BiMap<String, String> serverToEndpoints;
try {
serverToEndpoints =
pinotHelixResourceManager.getDataInstanceAdminEndpoints(serverToSegments.keySet());
@@ -226,6 +228,17 @@ public class UpsertCompactionTaskGenerator extends
BaseTaskGenerator {
segment.getCrc(), validDocIdsMetadata.getSegmentCrc());
continue;
}
+
+ // skipping segments for which their servers are not in READY state.
The bitmaps would be inconsistent when
+ // server is NOT READY as UPDATING segments might be updating the
ONLINE segments
+ if (validDocIdsMetadata.getServerStatus() != null &&
!validDocIdsMetadata.getServerStatus()
+ .equals(ServiceStatus.Status.GOOD)) {
+ LOGGER.warn("Server {} is in {} state, skipping {} generation for
segment: {}",
+ validDocIdsMetadata.getInstanceId(),
validDocIdsMetadata.getServerStatus(),
+ MinionConstants.UpsertCompactionTask.TASK_TYPE, segmentName);
+ continue;
+ }
+
long totalDocs = validDocIdsMetadata.getTotalDocs();
double invalidRecordPercent = ((double) totalInvalidDocs / totalDocs)
* 100;
if (totalInvalidDocs == totalDocs) {
@@ -234,15 +247,13 @@ public class UpsertCompactionTaskGenerator extends
BaseTaskGenerator {
} else if (invalidRecordPercent >= invalidRecordsThresholdPercent
&& totalInvalidDocs >= invalidRecordsThresholdCount) {
LOGGER.debug("Segment {} contains {} invalid records out of {} total
records "
- + "(count threshold: {}, percent threshold: {}),
adding it to the compaction list",
- segmentName, totalInvalidDocs, totalDocs,
invalidRecordsThresholdCount,
- invalidRecordsThresholdPercent);
+ + "(count threshold: {}, percent threshold: {}), adding it
to the compaction list", segmentName,
+ totalInvalidDocs, totalDocs, invalidRecordsThresholdCount,
invalidRecordsThresholdPercent);
segmentsForCompaction.add(Pair.of(segment, totalInvalidDocs));
} else {
LOGGER.debug("Segment {} contains {} invalid records out of {} total
records "
- + "(count threshold: {}, percent threshold: {}),
skipping it for compaction",
- segmentName, totalInvalidDocs, totalDocs,
invalidRecordsThresholdCount,
- invalidRecordsThresholdPercent);
+ + "(count threshold: {}, percent threshold: {}), skipping it
for compaction", segmentName,
+ totalInvalidDocs, totalDocs, invalidRecordsThresholdCount,
invalidRecordsThresholdPercent);
}
break;
}
diff --git
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskGenerator.java
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskGenerator.java
index 74dc529fc4..13097da140 100644
---
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskGenerator.java
+++
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/main/java/org/apache/pinot/plugin/minion/tasks/upsertcompactmerge/UpsertCompactMergeTaskGenerator.java
@@ -37,6 +37,7 @@ import
org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo;
import org.apache.pinot.common.restlet.resources.ValidDocIdsType;
import org.apache.pinot.common.utils.SegmentUtils;
+import org.apache.pinot.common.utils.ServiceStatus;
import org.apache.pinot.controller.helix.core.PinotHelixResourceManager;
import
org.apache.pinot.controller.helix.core.minion.generator.BaseTaskGenerator;
import
org.apache.pinot.controller.helix.core.minion.generator.TaskGeneratorUtils;
@@ -156,7 +157,8 @@ public class UpsertCompactMergeTaskGenerator extends
BaseTaskGenerator {
// get server to segment mappings
PinotHelixResourceManager pinotHelixResourceManager =
_clusterInfoAccessor.getPinotHelixResourceManager();
- Map<String, List<String>> serverToSegments =
pinotHelixResourceManager.getServerToSegmentsMap(tableNameWithType);
+ Map<String, List<String>> serverToSegments =
+
pinotHelixResourceManager.getServerToOnlineSegmentsMapFromEV(tableNameWithType,
true);
BiMap<String, String> serverToEndpoints;
try {
serverToEndpoints =
pinotHelixResourceManager.getDataInstanceAdminEndpoints(serverToSegments.keySet());
@@ -288,6 +290,16 @@ public class UpsertCompactMergeTaskGenerator extends
BaseTaskGenerator {
continue;
}
+ // skipping segments for which their servers are not in READY state.
The bitmaps would be inconsistent when
+ // server is NOT READY as UPDATING segments might be updating the
ONLINE segments
+ if (validDocIdsMetadata.getServerStatus() != null &&
!validDocIdsMetadata.getServerStatus()
+ .equals(ServiceStatus.Status.GOOD)) {
+ LOGGER.warn("Server {} is in {} state, skipping {} generation for
segment: {}",
+ validDocIdsMetadata.getInstanceId(),
validDocIdsMetadata.getServerStatus(),
+ MinionConstants.UpsertCompactMergeTask.TASK_TYPE, segmentName);
+ continue;
+ }
+
// segments eligible for deletion with no valid records
long totalDocs = validDocIdsMetadata.getTotalDocs();
if (totalInvalidDocs == totalDocs) {
diff --git
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java
index 5436487e54..410f8f5af2 100644
---
a/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java
+++
b/pinot-plugins/pinot-minion-tasks/pinot-minion-builtin-tasks/src/test/java/org/apache/pinot/plugin/minion/tasks/upsertcompaction/UpsertCompactionTaskGeneratorTest.java
@@ -222,10 +222,10 @@ public class UpsertCompactionTaskGeneratorTest {
Map<String, String> compactionConfigs = getCompactionConfigs("1", "10");
String json = "{\"testTable__0\": [{\"totalValidDocs\": 50,
\"totalInvalidDocs\": 50, "
+ "\"segmentName\": \"testTable__0\", \"totalDocs\": 100,
\"segmentCrc\": \"1000\", "
- + "\"segmentCreationTimeMillis\": 1234567890}], "
+ + "\"segmentCreationTimeMillis\": 1234567890, \"serverStatus\":
\"GOOD\", \"instanceId\": \"server1\"}], "
+ "\"testTable__1\": [{\"totalValidDocs\": 0, "
+ "\"totalInvalidDocs\": 10, \"segmentName\": \"testTable__1\",
\"totalDocs\": 10, \"segmentCrc\": \"2000\", "
- + "\"segmentCreationTimeMillis\": 9876543210}]}";
+ + "\"segmentCreationTimeMillis\": 9876543210, \"serverStatus\":
\"GOOD\", \"instanceId\": \"server1\"}]}";
Map<String, List<ValidDocIdsMetadataInfo>> validDocIdsMetadataInfo =
JsonUtils.stringToObject(json, new TypeReference<>() {
@@ -284,7 +284,8 @@ public class UpsertCompactionTaskGeneratorTest {
+ "\"1234567890\", \"segmentCreationTimeMillis\": 1111111111}], \"" +
_completedSegment2.getSegmentName()
+ "\": [{\"totalValidDocs\": 0, " + "\"totalInvalidDocs\": 10,
\"segmentName\": \""
+ _completedSegment2.getSegmentName() + "\", " + "\"segmentCrc\": \""
+ _completedSegment2.getCrc()
- + "\", \"totalDocs\": 10, \"segmentCreationTimeMillis\":
2222222222}]}";
+ + "\", \"totalDocs\": 10, \"segmentCreationTimeMillis\": 2222222222,
\"serverStatus\": \"GOOD\", "
+ + "\"instanceId\": \"server1\"}]}";
validDocIdsMetadataInfo = JsonUtils.stringToObject(json, new
TypeReference<>() {
});
segmentSelectionResult =
@@ -301,11 +302,14 @@ public class UpsertCompactionTaskGeneratorTest {
// check if both the candidates for compaction are coming in sorted
descending order
json = "{\"" + _completedSegment.getSegmentName() + "\":
[{\"totalValidDocs\": 50, \"totalInvalidDocs\": 50, "
+ "\"segmentName\": \"" + _completedSegment.getSegmentName() + "\",
\"totalDocs\": 100, \"segmentCrc\": \""
- + _completedSegment.getCrc() + "\", \"segmentCreationTimeMillis\":
1234567890}], \""
+ + _completedSegment.getCrc()
+ + "\", \"segmentCreationTimeMillis\": 1234567890, \"serverStatus\":
\"GOOD\", \"instanceId\": \"server1\"}]"
+ + ", \""
+ _completedSegment2.getSegmentName() + "\": "
+ "[{\"totalValidDocs\": 10, \"totalInvalidDocs\": 40,
\"segmentName\": \""
+ _completedSegment2.getSegmentName() + "\", \"segmentCrc\": \"" +
_completedSegment2.getCrc() + "\", "
- + "\"totalDocs\": 50, \"segmentCreationTimeMillis\": 9876543210}]}";
+ + "\"totalDocs\": 50, \"segmentCreationTimeMillis\": 9876543210,
\"serverStatus\": \"GOOD\", \"instanceId\": "
+ + "\"server1\"}]}";
validDocIdsMetadataInfo = JsonUtils.stringToObject(json, new
TypeReference<>() {
});
compactionConfigs = getCompactionConfigs("30", "0");
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 da8ba5732a..c93a70b1c5 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
@@ -79,6 +79,7 @@ import
org.apache.pinot.common.restlet.resources.ValidDocIdsType;
import org.apache.pinot.common.utils.DatabaseUtils;
import org.apache.pinot.common.utils.LLCSegmentName;
import org.apache.pinot.common.utils.RoaringBitmapUtils;
+import org.apache.pinot.common.utils.ServiceStatus;
import org.apache.pinot.common.utils.TarCompressionUtils;
import org.apache.pinot.common.utils.URIUtils;
import org.apache.pinot.common.utils.helix.HelixHelper;
@@ -104,6 +105,7 @@ import org.apache.pinot.spi.config.table.TableType;
import org.apache.pinot.spi.data.FieldSpec;
import org.apache.pinot.spi.data.FieldSpec.DataType;
import org.apache.pinot.spi.stream.ConsumerPartitionState;
+import org.apache.pinot.spi.utils.CommonConstants;
import
org.apache.pinot.spi.utils.CommonConstants.Helix.StateModel.SegmentStateModel;
import org.apache.pinot.spi.utils.JsonUtils;
import org.apache.pinot.spi.utils.builder.TableNameBuilder;
@@ -111,19 +113,18 @@ import org.roaringbitmap.buffer.MutableRoaringBitmap;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import static org.apache.pinot.spi.utils.CommonConstants.DATABASE;
-import static
org.apache.pinot.spi.utils.CommonConstants.SWAGGER_AUTHORIZATION_KEY;
-
-@Api(tags = "Table", authorizations = {@Authorization(value =
SWAGGER_AUTHORIZATION_KEY),
- @Authorization(value = DATABASE)})
+@Api(tags = "Table", authorizations = {
+ @Authorization(value = CommonConstants.SWAGGER_AUTHORIZATION_KEY),
@Authorization(value = CommonConstants.DATABASE)
+})
@SwaggerDefinition(securityDefinition =
@SecurityDefinition(apiKeyAuthDefinitions = {
- @ApiKeyAuthDefinition(name = HttpHeaders.AUTHORIZATION, in =
ApiKeyAuthDefinition.ApiKeyLocation.HEADER,
- key = SWAGGER_AUTHORIZATION_KEY,
- description = "The format of the key is ```\"Basic <token>\" or
\"Bearer <token>\"```"),
- @ApiKeyAuthDefinition(name = DATABASE, in =
ApiKeyAuthDefinition.ApiKeyLocation.HEADER, key = DATABASE,
- description = "Database context passed through http header. If no
context is provided 'default' database "
- + "context will be considered.")}))
+ @ApiKeyAuthDefinition(name = HttpHeaders.AUTHORIZATION, in =
ApiKeyAuthDefinition.ApiKeyLocation.HEADER, key =
+ CommonConstants.SWAGGER_AUTHORIZATION_KEY, description = "The format
of the key is ```\"Basic <token>\" or "
+ + "\"Bearer <token>\"```"), @ApiKeyAuthDefinition(name =
CommonConstants.DATABASE, in =
+ ApiKeyAuthDefinition.ApiKeyLocation.HEADER, key =
CommonConstants.DATABASE, description =
+ "Database context passed through http header. If no context is provided
'default' database "
+ + "context will be considered.")
+}))
@Path("/")
public class TablesResource {
private static final Logger LOGGER =
LoggerFactory.getLogger(TablesResource.class);
@@ -502,8 +503,7 @@ public class TablesResource {
public ValidDocIdsBitmapResponse downloadValidDocIdsBitmap(
@ApiParam(value = "Name of the table with type REALTIME", required =
true, example = "myTable_REALTIME")
@PathParam("tableNameWithType") String tableNameWithType,
- @ApiParam(value = "Valid doc ids type")
- @QueryParam("validDocIdsType") String validDocIdsType,
+ @ApiParam(value = "Valid doc ids type") @QueryParam("validDocIdsType")
String validDocIdsType,
@ApiParam(value = "Name of the segment", required = true)
@PathParam("segmentName") @Encoded String segmentName,
@Context HttpHeaders httpHeaders) {
tableNameWithType = DatabaseUtils.translateTableName(tableNameWithType,
httpHeaders);
@@ -528,6 +528,7 @@ public class TablesResource {
String.format("Table %s segment %s is not a immutable segment",
tableNameWithType, segmentName),
Response.Status.BAD_REQUEST);
}
+ ServiceStatus.Status status =
ServiceStatus.getServiceStatus(_instanceId);
final Pair<ValidDocIdsType, MutableRoaringBitmap>
validDocIdsSnapshotPair =
getValidDocIds(indexSegment, validDocIdsType);
@@ -537,15 +538,15 @@ public class TablesResource {
if (validDocIdSnapshot == null) {
String msg = String.format(
"Found that validDocIds is missing while fetching validDocIds for
table %s segment %s while "
- + "reading the validDocIds with validDocIdType %s",
- tableNameWithType, segmentDataManager.getSegmentName(),
validDocIdsType);
+ + "reading the validDocIds with validDocIdType %s",
tableNameWithType,
+ segmentDataManager.getSegmentName(), validDocIdsType);
LOGGER.warn(msg);
throw new WebApplicationException(msg, Response.Status.NOT_FOUND);
}
-
byte[] validDocIdsBytes =
RoaringBitmapUtils.serialize(validDocIdSnapshot);
return new ValidDocIdsBitmapResponse(segmentName,
indexSegment.getSegmentMetadata().getCrc(),
- finalValidDocIdsType, validDocIdsBytes);
+ finalValidDocIdsType, validDocIdsBytes,
_serverInstance.getInstanceDataManager().getInstanceId(),
+ status);
} finally {
tableDataManager.releaseSegment(segmentDataManager);
}
@@ -627,7 +628,8 @@ public class TablesResource {
@ApiParam(value = "Valid doc ids type")
@QueryParam("validDocIdsType") String validDocIdsType,
@ApiParam(value = "Segment name", allowMultiple = true)
@QueryParam("segmentNames") List<String> segmentNames,
- @Context HttpHeaders headers) {
+ @Context HttpHeaders headers)
+ throws Exception {
tableNameWithType = DatabaseUtils.translateTableName(tableNameWithType,
headers);
return ResourceUtils.convertToJsonString(
processValidDocIdsMetadata(tableNameWithType, segmentNames,
validDocIdsType));
@@ -645,9 +647,8 @@ public class TablesResource {
public String getValidDocIdsMetadata(
@ApiParam(value = "Table name including type", required = true, example
= "myTable_REALTIME")
@PathParam("tableNameWithType") String tableNameWithType,
- @ApiParam(value = "Valid doc ids type")
- @QueryParam("validDocIdsType") String validDocIdsType, TableSegments
tableSegments,
- @Context HttpHeaders headers) {
+ @ApiParam(value = "Valid doc ids type") @QueryParam("validDocIdsType")
String validDocIdsType,
+ TableSegments tableSegments, @Context HttpHeaders headers) {
tableNameWithType = DatabaseUtils.translateTableName(tableNameWithType,
headers);
List<String> segmentNames = tableSegments.getSegments();
return ResourceUtils.convertToJsonString(
@@ -660,7 +661,7 @@ public class TablesResource {
ServerResourceUtils.checkGetTableDataManager(_serverInstance,
tableNameWithType);
List<String> missingSegments = new ArrayList<>();
int nonImmutableSegmentCount = 0;
- int missingValidDocIdSnapshotSegmentCount = 0;
+ List<String> missingValidDocsSegments = new ArrayList<>();
List<SegmentDataManager> segmentDataManagers;
if (segments == null) {
segmentDataManagers = tableDataManager.acquireAllSegments();
@@ -675,6 +676,8 @@ public class TablesResource {
// process the remaining available segments.
LOGGER.warn("Table {} has missing segments {}", tableNameWithType,
missingSegments);
}
+ ServiceStatus.Status status =
ServiceStatus.getServiceStatus(_instanceId);
+
List<Map<String, Object>> allValidDocIdsMetadata = new
ArrayList<>(segmentDataManagers.size());
for (SegmentDataManager segmentDataManager : segmentDataManagers) {
IndexSegment indexSegment = segmentDataManager.getSegment();
@@ -705,7 +708,7 @@ public class TablesResource {
segmentDataManager.getSegmentName(), validDocIdsType);
LOGGER.debug(msg);
}
- missingValidDocIdSnapshotSegmentCount++;
+ missingValidDocsSegments.add(segmentDataManager.getSegmentName());
continue;
}
@@ -719,6 +722,8 @@ public class TablesResource {
validDocIdsMetadata.put("totalInvalidDocs", totalInvalidDocs);
validDocIdsMetadata.put("segmentCrc",
indexSegment.getSegmentMetadata().getCrc());
validDocIdsMetadata.put("validDocIdsType", finalValidDocIdsType);
+ validDocIdsMetadata.put("serverStatus", status);
+ validDocIdsMetadata.put("instanceId",
_serverInstance.getInstanceDataManager().getInstanceId());
if (segmentDataManager instanceof ImmutableSegmentDataManager) {
validDocIdsMetadata.put("segmentSizeInBytes",
((ImmutableSegment)
segmentDataManager.getSegment()).getSegmentSizeBytes());
@@ -730,10 +735,10 @@ public class TablesResource {
LOGGER.warn("Table {} has {} non-immutable segments found while
processing validDocIdsMetadata",
tableNameWithType, nonImmutableSegmentCount);
}
- if (missingValidDocIdSnapshotSegmentCount > 0) {
- LOGGER.warn("Found that validDocIds is missing for {} segments while
processing validDocIdsMetadata "
- + "for table {} while reading the validDocIds with
validDocIdType {}. ",
- missingValidDocIdSnapshotSegmentCount, tableNameWithType,
validDocIdsType);
+ if (!missingValidDocsSegments.isEmpty()) {
+ LOGGER.warn("Found that validDocIds is missing for segments {} while
processing validDocIdsMetadata "
+ + "for table {} while reading the validDocIds with
validDocIdType {}. ", missingValidDocsSegments,
+ tableNameWithType, validDocIdsType);
}
return allValidDocIdsMetadata;
} finally {
diff --git
a/pinot-server/src/test/java/org/apache/pinot/server/api/BaseResourceTest.java
b/pinot-server/src/test/java/org/apache/pinot/server/api/BaseResourceTest.java
index 9967cade27..93b6191114 100644
---
a/pinot-server/src/test/java/org/apache/pinot/server/api/BaseResourceTest.java
+++
b/pinot-server/src/test/java/org/apache/pinot/server/api/BaseResourceTest.java
@@ -30,6 +30,7 @@ import java.util.concurrent.Executors;
import javax.ws.rs.client.ClientBuilder;
import javax.ws.rs.client.WebTarget;
import org.apache.commons.io.FileUtils;
+import org.apache.helix.HelixAdmin;
import org.apache.helix.HelixManager;
import org.apache.pinot.common.config.TlsConfig;
import org.apache.pinot.common.metrics.ServerMetrics;
@@ -93,6 +94,7 @@ public abstract class BaseResourceTest {
protected AdminApiApplication _adminApiApplication;
protected WebTarget _webTarget;
protected String _instanceId;
+ protected ServerInstance _serverInstance;
@SuppressWarnings("SuspiciousMethodCalls")
@BeforeClass
@@ -113,12 +115,19 @@ public abstract class BaseResourceTest {
when(instanceDataManager.getAllTables()).thenReturn(_tableDataManagerMap.keySet());
// Mock the server instance
- ServerInstance serverInstance = mock(ServerInstance.class);
-
when(serverInstance.getServerMetrics()).thenReturn(mock(ServerMetrics.class));
-
when(serverInstance.getInstanceDataManager()).thenReturn(instanceDataManager);
-
when(serverInstance.getInstanceDataManager().getSegmentFileDirectory()).thenReturn(
+ _serverInstance = mock(ServerInstance.class);
+
when(_serverInstance.getServerMetrics()).thenReturn(mock(ServerMetrics.class));
+
when(_serverInstance.getInstanceDataManager()).thenReturn(instanceDataManager);
+
when(_serverInstance.getInstanceDataManager().getSegmentFileDirectory()).thenReturn(
FileUtils.getTempDirectoryPath());
-
when(serverInstance.getHelixManager()).thenReturn(mock(HelixManager.class));
+
+ // Create a single HelixManager mock with proper segment data
+ HelixManager helixManager = mock(HelixManager.class);
+ HelixAdmin helixAdmin = mock(HelixAdmin.class);
+ when(helixManager.getClusterManagmentTool()).thenReturn(helixAdmin);
+ when(helixManager.getClusterName()).thenReturn("testCluster");
+
+ when(_serverInstance.getHelixManager()).thenReturn(helixManager);
// Mock the segment uploader
SegmentUploader segmentUploader = mock(SegmentUploader.class);
@@ -145,7 +154,7 @@ public abstract class BaseResourceTest {
CommonConstants.Helix.DEFAULT_SERVER_NETTY_PORT);
_instanceId = CommonConstants.Helix.PREFIX_OF_SERVER_INSTANCE + hostname +
"_" + port;
serverConf.setProperty(CommonConstants.Server.CONFIG_OF_INSTANCE_ID,
_instanceId);
- _adminApiApplication = new AdminApiApplication(serverInstance, new
AllowAllAccessFactory(), serverConf);
+ _adminApiApplication = new AdminApiApplication(_serverInstance, new
AllowAllAccessFactory(), serverConf);
_adminApiApplication.start(Collections.singletonList(
new ListenerConfig(CommonConstants.HTTP_PROTOCOL, "0.0.0.0",
CommonConstants.Server.DEFAULT_ADMIN_API_PORT,
CommonConstants.HTTP_PROTOCOL, new TlsConfig(),
HttpServerThreadPoolConfig.defaultInstance())));
@@ -200,15 +209,20 @@ public abstract class BaseResourceTest {
protected void addTable(String tableNameWithType) {
InstanceDataManagerConfig instanceDataManagerConfig =
mock(InstanceDataManagerConfig.class);
when(instanceDataManagerConfig.getInstanceDataDir()).thenReturn(TEMP_DIR.getAbsolutePath());
+
when(instanceDataManagerConfig.getInstanceId()).thenReturn("Server_1_100.89.121.12");
TableType tableType =
TableNameBuilder.getTableTypeFromTableName(tableNameWithType);
assertNotNull(tableType);
TableConfig tableConfig = new
TableConfigBuilder(tableType).setTableName(tableNameWithType).build();
Schema schema =
new
Schema.SchemaBuilder().setSchemaName(TableNameBuilder.extractRawTableName(tableNameWithType)).build();
+
+ // Get the HelixManager from the server instance (already configured in
setUp)
+ HelixManager helixManager = _serverInstance.getHelixManager();
+
// NOTE: Use OfflineTableDataManager for both OFFLINE and REALTIME table
because RealtimeTableDataManager performs
// more checks
TableDataManager tableDataManager = new OfflineTableDataManager();
- tableDataManager.init(instanceDataManagerConfig, mock(HelixManager.class),
new SegmentLocks(), tableConfig, schema,
+ tableDataManager.init(instanceDataManagerConfig, helixManager, new
SegmentLocks(), tableConfig, schema,
new SegmentReloadSemaphore(1), Executors.newSingleThreadExecutor(),
null, null, null);
tableDataManager.start();
_tableDataManagerMap.put(tableNameWithType, tableDataManager);
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 ffd06c8b42..2ed2af463d 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
@@ -354,6 +354,12 @@ public class TablesResourceTest extends BaseResourceTest {
Assert.assertEquals(validDocIdsMetadata.get("segmentSizeInBytes").asLong(),
1877636);
Assert.assertTrue(validDocIdsMetadata.has("segmentCreationTimeMillis"));
Assert.assertTrue(validDocIdsMetadata.get("segmentCreationTimeMillis").asLong()
> 0);
+
+ // Verify server status information
+ Assert.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);
}
// Verify metadata file from segments.
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]