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]

Reply via email to