goutamadwant commented on code in PR #19178:
URL: https://github.com/apache/pinot/pull/19178#discussion_r4117214738
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -1351,6 +1453,13 @@ public Set<String> getServingInstances(String
tableNameWithType) {
return routingEntry._instanceSelector.getServingInstances();
}
+ /// Returns whether the broker sees the server as routable.
+ public boolean isServerRoutable(String instanceId) {
+ // Read the map first. A new entry is inserted only after the server is
marked pending, and removing the pending
+ // marker publishes all routing-entry updates that precede it.
+ return _routableServerInstanceMap.containsKey(instanceId) &&
!_serversPendingRoutingUpdate.contains(instanceId);
Review Comment:
Addressed. Re-included servers are now marked pending before publication and
remain unready until all relevant selector updates succeed; failed updates are
retained and retried.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -486,16 +520,29 @@ private void processInstanceConfigChangeInternal() {
updateRoutingEntryOnInstancesChange(routingEntry,
_routableServerInstanceMap.keySet(), changedServers);
}
} catch (Exception e) {
- LOGGER.error("Caught unexpected exception while updating routing entry
on instances change for table: {}",
- tableNameWithType, e);
+ routingUpdateSucceeded = false;
+ logInstanceConfigRetryFailure(
+ "Caught unexpected exception while updating routing entry on
instances change for table: "
+ + tableNameWithType, e);
}
}
long updateRoutingEntriesEndTimeMs = System.currentTimeMillis();
+ if (routingUpdateSucceeded) {
Review Comment:
Addressed. Readiness is now tracked per server against only the tables that
assign that server, so an unrelated table failure no longer blocks it.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -450,14 +475,21 @@ private void processInstanceConfigChangeInternal() {
}
// Calculate the routable servers and the changed routable servers
- List<String> changedServers = new ArrayList<>(newEnabledServers.size() +
newDisabledServers.size());
+ // Include servers left pending by an aborted prior refresh so the next
callback retries their routing publication.
+ List<String> pendingEnabledServers = new ArrayList<>();
+ for (String server : enabledServers) {
+ if (_serversPendingRoutingUpdate.contains(server)) {
+ pendingEnabledServers.add(server);
+ }
+ }
+ List<String> changedServers = new ArrayList<>(pendingEnabledServers.size()
+ newDisabledServers.size());
if (_excludedServers.isEmpty()) {
- changedServers.addAll(newEnabledServers);
+ changedServers.addAll(pendingEnabledServers);
changedServers.addAll(newDisabledServers);
} else {
enabledServers.removeAll(_excludedServers);
- // NOTE: All new enabled servers are routable
- changedServers.addAll(newEnabledServers);
+ // NOTE: All pending enabled servers are routable
+ changedServers.addAll(pendingEnabledServers);
Review Comment:
Addressed. Newly disabled servers are added to the changed set independently
of exclusions, with regression coverage for excluding one server while
disabling another.
##########
pinot-common/src/main/java/org/apache/pinot/common/utils/http/HttpClient.java:
##########
@@ -286,23 +331,33 @@ public SimpleHttpResponse sendRequest(ClassicHttpRequest
request, long socketTim
RequestConfig requestConfig =
RequestConfig.custom()
.setResponseTimeout(Timeout.ofMilliseconds(socketTimeoutMs))
-
.setConnectionRequestTimeout(Timeout.ofMilliseconds(DEFAULT_CONNECTION_REQUEST_TIMEOUT_MS))
+
.setConnectionRequestTimeout(Timeout.ofMilliseconds(connectionRequestTimeoutMs))
.build();
HttpClientContext clientContext = HttpClientContext.create();
clientContext.setRequestConfig(requestConfig);
- try (CloseableHttpResponse response = _httpClient.execute(request,
clientContext)) {
+ CloseableHttpResponse response = _httpClient.execute(request,
clientContext);
+ boolean closeImmediately = false;
+ try {
if
(response.containsHeader(CommonConstants.Controller.HOST_HTTP_HEADER)) {
String controllerHost =
response.getFirstHeader(CommonConstants.Controller.HOST_HTTP_HEADER).getValue();
String controllerVersion =
response.getFirstHeader(CommonConstants.Controller.VERSION_HTTP_HEADER).getValue();
LOGGER.info("Sending request: {} to controller: {}, version: {}",
request.getRequestUri(), controllerHost,
controllerVersion);
}
int statusCode = response.getCode();
+ BoundedResponseContent responseContent =
readResponseContent(response.getEntity(), maxResponseLength);
Review Comment:
Addressed. Error-body read failures now preserve the HTTP status, return the
fallback reason, and close the response immediately.
##########
pinot-common/src/main/java/org/apache/pinot/common/utils/http/HttpClient.java:
##########
@@ -555,10 +655,29 @@ private static CloseableHttpClient
buildCloseableHttpClient(HttpClientConfig htt
if (httpClientConfig.isDisableDefaultUserAgent()) {
httpClientBuilder.disableDefaultUserAgent();
}
+ if (!httpClientConfig.isFollowRedirects()) {
+ httpClientBuilder.disableRedirectHandling();
+ }
return httpClientBuilder.build();
}
private static String getErrorMessage(ClassicHttpRequest request,
CloseableHttpResponse response) {
+ return getErrorMessage(request, response, Integer.MAX_VALUE);
+ }
+
+ private static String getErrorMessage(ClassicHttpRequest request,
CloseableHttpResponse response,
Review Comment:
Addressed. Removed the unused max-response-length overloads and folded the
unbounded behavior into the existing wrappers.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]