Jackie-Jiang commented on code in PR #19178:
URL: https://github.com/apache/pinot/pull/19178#discussion_r4113607380
##########
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:
[MAJOR] This can return true before a re-included server is actually back in
table routing. includeServerToRoutingInternal publishes the routable map before
calling each selector, does not mark the server pending, and catches
selector-update failures. The endpoint can therefore return 200 during the
update or indefinitely after a failure; the server's startup checker then
caches READY. Re-inclusions need a pending marker until all selector updates
succeed, with a retry path on failure.
##########
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:
[MINOR] This reads an error response body before handling its status. If the
entity stream throws while being read, sendRequest now propagates the
exception; previously getErrorMessage caught body-read errors and returned a
SimpleHttpResponse containing the HTTP status and a fallback reason. Preserving
that fallback for status >= 300 would avoid changing existing callers' handling
of malformed or interrupted error bodies.
##########
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:
[MAJOR] routingUpdateSucceeded is global across every routing entry, so one
failing table keeps every newly enabled server pending. A server whose own
tables have already updated can return 503 from the routing endpoint solely
because an unrelated table's selector or health-metric update keeps failing.
With failOpen=false, that healthy server remains unready indefinitely. Could
the pending acknowledgement be scoped to tables that can route each server,
rather than all tables?
##########
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:
[MAJOR, preexisting] This branch has an existing disable-path bug: if server
X is excluded and a different server Y becomes disabled, Y is omitted from
changedServers because the loop below only adds disabled servers already in
_excludedServers. With no pending additions, the early return skips both
selector updates and removal of Y from _enabledServerInstanceMap, so the broker
may continue routing to disabled Y. The branch can be simplified by removing
exclusions from enabledServers, then adding all pending and all newly disabled
IDs to changedServers unconditionally.
##########
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:
[Simplification] This int maxResponseLength overload has no bounded caller;
only the two-argument wrapper invokes it with Integer.MAX_VALUE. The new
httpEntityToString(HttpEntity, int) overload is likewise called only by its
unbounded wrapper. Folding both bodies into their wrappers removes unused
generality while keeping the String error-message overload needed by bounded
sendRequest.
--
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]