Jackie-Jiang commented on code in PR #19178:
URL: https://github.com/apache/pinot/pull/19178#discussion_r4187568728
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -449,26 +479,28 @@ private void processInstanceConfigChangeInternal() {
}
}
- // Calculate the routable servers and the changed routable servers
- List<String> changedServers = new ArrayList<>(newEnabledServers.size() +
newDisabledServers.size());
- if (_excludedServers.isEmpty()) {
- changedServers.addAll(newEnabledServers);
- changedServers.addAll(newDisabledServers);
- } else {
- enabledServers.removeAll(_excludedServers);
- // NOTE: All new enabled servers are routable
- changedServers.addAll(newEnabledServers);
- for (String newDisabledServer : newDisabledServers) {
- if (_excludedServers.contains(newDisabledServer)) {
- changedServers.add(newDisabledServer);
- }
+ // Calculate the routable servers and the changed routable servers.
Include servers left pending by an aborted prior
+ // refresh so the next callback retries their routing publication.
+ List<String> pendingEnabledServers = new ArrayList<>();
Review Comment:
[Performance] On every retry, a server still pending acknowledgement is
added to `changedServers`, which is then passed to every table. One
persistently failing table therefore makes already successful selectors rebuild
repeatedly under the global write lock. Could retries process only
`_pendingRoutingUpdatesByTable` after the initial pass? A callback that aborts
before its first table pass still needs explicit recovery.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -477,25 +509,23 @@ private void processInstanceConfigChangeInternal() {
return;
}
- // Update routing entry for all tables
- for (RoutingEntry routingEntry : _routingEntryMap.values()) {
- String tableNameWithType = routingEntry.getTableNameWithType();
- try {
- Object tableLock = getRoutingTableBuildLock(tableNameWithType);
- synchronized (tableLock) {
- updateRoutingEntryOnInstancesChange(routingEntry,
_routableServerInstanceMap.keySet(), changedServers);
- }
- } catch (Exception e) {
- LOGGER.error("Caught unexpected exception while updating routing entry
on instances change for table: {}",
- tableNameWithType, e);
- }
- }
+ updateRoutingEntriesOnInstancesChange(changedServers);
long updateRoutingEntriesEndTimeMs = System.currentTimeMillis();
+ if (_pendingRoutingUpdatesByTable.isEmpty()) {
Review Comment:
[Simplification] The same reset-delay-or-schedule-retry decision appears in
the instance-config, include, and exclude paths. A small helper for the
post-update retry decision would keep these paths consistent and remove the
repeated block.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -449,26 +479,28 @@ private void processInstanceConfigChangeInternal() {
}
}
- // Calculate the routable servers and the changed routable servers
- List<String> changedServers = new ArrayList<>(newEnabledServers.size() +
newDisabledServers.size());
- if (_excludedServers.isEmpty()) {
- changedServers.addAll(newEnabledServers);
- changedServers.addAll(newDisabledServers);
- } else {
- enabledServers.removeAll(_excludedServers);
- // NOTE: All new enabled servers are routable
- changedServers.addAll(newEnabledServers);
- for (String newDisabledServer : newDisabledServers) {
- if (_excludedServers.contains(newDisabledServer)) {
- changedServers.add(newDisabledServer);
- }
+ // Calculate the routable servers and the changed routable servers.
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.containsKey(server)) {
+ pendingEnabledServers.add(server);
}
}
- _routableServerInstanceMap =
buildRoutableServerInstanceMap(enabledServers);
+ List<String> changedServers = new ArrayList<>(pendingEnabledServers.size()
+ newDisabledServers.size());
+ changedServers.addAll(pendingEnabledServers);
+ // A disabled server must reach every selector regardless of whether a
different server is excluded. Otherwise the
+ // disabled server can remain in the selectors indefinitely.
+ changedServers.addAll(newDisabledServers);
+ markServersPendingRelevantTableUpdates(pendingEnabledServers);
Review Comment:
[Performance / rollout] The server-side readiness gate defaults off, but
this broker still scans every routing entry for each newly enabled or pending
server, maintains acknowledgements, and can schedule retries. That may be
deliberate so brokers are ready before the server flag is enabled; if so, it
would help to document the feature-off broker cost. Otherwise a broker-side
switch would be needed for zero tracking cost while off. The two hotspots below
are worth reducing either way.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java:
##########
@@ -507,6 +537,138 @@ private void processInstanceConfigChangeInternal() {
newDisabledServers, _excludedServers);
}
+ /// Replaces each server's publication marker with the exact tables whose
primary selector can route to it. Must run
+ /// before the server is published in `_routableServerInstanceMap` so the
lock-free readiness check cannot observe a
+ /// routable server without also observing its pending acknowledgements.
+ @GuardedBy("_globalLock.writeLock()")
+ private void markServersPendingRelevantTableUpdates(List<String> servers) {
+ for (String server : servers) {
+ Set<String> relevantTables = new HashSet<>();
+ for (RoutingEntry routingEntry : _routingEntryMap.values()) {
+ if (routingEntry.isServerAssigned(server)) {
+ relevantTables.add(routingEntry.getTableNameWithType());
+ }
+ }
+ _serversPendingRoutingUpdate.put(server, Set.copyOf(relevantTables));
+ }
+ }
+
+ /// Applies the current instance change plus any retained failed work
independently to each table. Successful tables
+ /// acknowledge only the servers they can route; failed tables retain their
complete change set for the retry.
+ @GuardedBy("_globalLock.writeLock()")
+ private void updateRoutingEntriesOnInstancesChange(List<String>
changedServers) {
+ for (RoutingEntry routingEntry : _routingEntryMap.values()) {
+ String tableNameWithType = routingEntry.getTableNameWithType();
+ List<String> tableChangedServers = new ArrayList<>(changedServers);
+ Set<String> pendingServers =
_pendingRoutingUpdatesByTable.get(tableNameWithType);
+ if (pendingServers != null) {
+ for (String pendingServer : pendingServers) {
+ if (!tableChangedServers.contains(pendingServer)) {
+ tableChangedServers.add(pendingServer);
+ }
+ }
+ }
+ if (tableChangedServers.isEmpty()) {
+ continue;
+ }
+
+ try {
+ Object tableLock = getRoutingTableBuildLock(tableNameWithType);
+ synchronized (tableLock) {
+ updateRoutingEntryOnInstancesChange(routingEntry,
_routableServerInstanceMap.keySet(), tableChangedServers);
+ }
+ _pendingRoutingUpdatesByTable.remove(tableNameWithType);
+ acknowledgeRoutingUpdate(tableNameWithType);
+ } catch (Exception e) {
+ _pendingRoutingUpdatesByTable.put(tableNameWithType,
Set.copyOf(tableChangedServers));
+ logInstanceConfigRetryFailure(
+ "Caught unexpected exception while updating routing entry on
instances change for table: "
+ + tableNameWithType, e);
+ }
+ }
+
+ // A table can disappear after work was retained but before a retry. It
can no longer route to any server, so it
+ // must neither be retried forever nor keep a readiness acknowledgement
pending.
+ for (String tableNameWithType :
Set.copyOf(_pendingRoutingUpdatesByTable.keySet())) {
+ if (!_routingEntryMap.containsKey(tableNameWithType)) {
+ _pendingRoutingUpdatesByTable.remove(tableNameWithType);
+ acknowledgeRoutingUpdate(tableNameWithType);
+ }
+ }
+
+ // Empty values are publication markers for servers without a relevant
table. They can be acknowledged only after
+ // the routable map has been published and the update pass has completed.
+ _serversPendingRoutingUpdate.entrySet().removeIf(entry ->
entry.getValue().isEmpty());
+ }
+
+ private void acknowledgeRoutingUpdate(String tableNameWithType) {
+ for (String server : Set.copyOf(_serversPendingRoutingUpdate.keySet())) {
+ _serversPendingRoutingUpdate.computeIfPresent(server, (key,
pendingTables) -> {
+ if (!pendingTables.contains(tableNameWithType)) {
+ return pendingTables;
+ }
+ Set<String> remainingTables = new HashSet<>(pendingTables);
Review Comment:
[Performance / simplification] Each table acknowledgement copies a server's
remaining table set into a `HashSet` and then into `Set.copyOf`. For a server
assigned to T tables, this copies O(T^2) table entries during enablement. The
readiness path only checks map membership, so a mutable per-server set updated
inside `computeIfPresent` would avoid those copies while serializing concurrent
acknowledgements. The outer `Set.copyOf(keySet())` also appears unnecessary;
direct iteration over the `ConcurrentHashMap` key set handles concurrent
removals.
##########
pinot-common/src/main/java/org/apache/pinot/common/utils/http/HttpClient.java:
##########
@@ -277,6 +314,21 @@ public SimpleHttpResponse sendRequest(ClassicHttpRequest
request)
public SimpleHttpResponse sendRequest(ClassicHttpRequest request, long
socketTimeoutMs)
throws IOException {
+ return sendRequest(request, socketTimeoutMs,
DEFAULT_CONNECTION_REQUEST_TIMEOUT_MS);
+ }
+
+ public SimpleHttpResponse sendRequest(ClassicHttpRequest request, long
socketTimeoutMs,
+ long connectionRequestTimeoutMs)
+ throws IOException {
+ return sendRequest(request, socketTimeoutMs, connectionRequestTimeoutMs,
Integer.MAX_VALUE);
+ }
+
+ // The response-handler overloads consume the complete entity before
returning. This raw-response overload is
+ // required so bounded reads can close a truncated response immediately
without draining the remaining body.
+ @SuppressWarnings("deprecation")
+ public SimpleHttpResponse sendRequest(ClassicHttpRequest request, long
socketTimeoutMs,
+ long connectionRequestTimeoutMs, int maxResponseLength)
Review Comment:
[API simplification] The new public `maxResponseLength` argument is checked
only while reading the response. A negative value throws on 2xx but is caught
as an error-body failure for status >= 300, so the same invalid argument
behaves differently by status. Validating `maxResponseLength >= 0` at this
method's entry would make the contract consistent; current production callers
use positive limits. Zero can remain a valid status-only bound, with an
optional fast path that avoids opening the entity stream.
--
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]