Jackie-Jiang commented on code in PR #19178:
URL: https://github.com/apache/pinot/pull/19178#discussion_r4110571675
##########
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 enabled.
+ public boolean isServerEnabled(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 _enabledServerInstanceMap.containsKey(instanceId) &&
!_serversPendingRoutingUpdate.contains(instanceId);
Review Comment:
[MAJOR] This returns true for a server that the failure detector has
excluded. Exclusion removes it from _routableServerInstanceMap but leaves
_enabledServerInstanceMap intact; the new test confirms that behavior. After an
abrupt crash and restart, the endpoint can acknowledge the server while the
broker still routes no queries to it. Using the existing routable snapshot
together with the per-server pending marker would match the PR's readiness
contract.
##########
pinot-server/src/main/java/org/apache/pinot/server/starter/helix/BrokerRoutingReadyChecker.java:
##########
@@ -0,0 +1,375 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.server.starter.helix;
+
+import com.google.common.annotations.VisibleForTesting;
+import com.google.common.util.concurrent.ThreadFactoryBuilder;
+import java.io.IOException;
+import java.net.URI;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.BooleanSupplier;
+import java.util.function.LongSupplier;
+import java.util.function.Predicate;
+import java.util.function.Supplier;
+import javax.annotation.Nullable;
+import org.apache.hc.core5.http.Header;
+import org.apache.helix.HelixAdmin;
+import org.apache.helix.HelixManager;
+import org.apache.helix.model.ExternalView;
+import org.apache.helix.model.InstanceConfig;
+import org.apache.pinot.common.auth.AuthProviderUtils;
+import org.apache.pinot.common.auth.NullAuthProvider;
+import org.apache.pinot.common.utils.SimpleHttpResponse;
+import org.apache.pinot.common.utils.config.InstanceUtils;
+import org.apache.pinot.common.utils.helix.HelixHelper;
+import org.apache.pinot.common.utils.http.HttpClient;
+import org.apache.pinot.common.utils.http.HttpClientConfig;
+import org.apache.pinot.common.utils.tls.TlsUtils;
+import org.apache.pinot.spi.auth.AuthProvider;
+import org.apache.pinot.spi.utils.CommonConstants;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+
+/// Checks broker routing state while a server starts. A single background
thread polls every online broker until all
+/// of them report the server as routable, then caches success. Readiness
requests only read the cached state and never
+/// perform network I/O. Broker requests are sequential and have bounded
connect, pool-checkout and response timeouts,
+/// so an unavailable broker cannot create unbounded tasks or threads.
+public class BrokerRoutingReadyChecker implements AutoCloseable {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(BrokerRoutingReadyChecker.class);
+ private static final long CHECK_INTERVAL_MS = 1_000L;
+ private static final long REQUEST_TIMEOUT_MS = 5_000L;
+ private static final int MAX_RESPONSE_LENGTH = 1_024;
+
+ private final String _serverInstanceId;
+ private final Supplier<Set<String>> _onlineBrokersSupplier;
+ private final BrokersReadyEvaluator _allBrokersReady;
+ @Nullable
+ private final ScheduledExecutorService _checkExecutor;
+ private final RoutingStatusClient _routingStatusClient;
+ private final LongSupplier _currentTimeMs;
+ private final long _deadlineMs;
+ private final boolean _failOpen;
+ private final AuthProvider _authProvider;
+ private final AtomicReference<State> _state;
+ private boolean _timeoutLogged;
+
+ public BrokerRoutingReadyChecker(HelixManager helixManager, long timeoutMs,
boolean failOpen,
+ AuthProvider authProvider) {
+ this(createProductionContext(helixManager, authProvider), timeoutMs,
failOpen, authProvider);
+ }
+
+ @VisibleForTesting
+ BrokerRoutingReadyChecker(HelixManager helixManager, long timeoutMs, boolean
failOpen, AuthProvider authProvider,
+ RoutingStatusClient routingStatusClient) {
+ this(helixManager, timeoutMs, failOpen, authProvider, routingStatusClient,
System::currentTimeMillis);
+ }
+
+ @VisibleForTesting
+ BrokerRoutingReadyChecker(HelixManager helixManager, long timeoutMs, boolean
failOpen, AuthProvider authProvider,
+ RoutingStatusClient routingStatusClient, LongSupplier currentTimeMs) {
+ this(createContext(helixManager, routingStatusClient, null, authProvider,
new AtomicReference<>(State.CHECKING)),
+ timeoutMs, failOpen, currentTimeMs, authProvider);
+ }
+
+ private BrokerRoutingReadyChecker(ProductionContext context, long timeoutMs,
boolean failOpen,
+ AuthProvider authProvider) {
+ this(context, timeoutMs, failOpen, System::currentTimeMillis,
authProvider);
+ }
+
+ private BrokerRoutingReadyChecker(ProductionContext context, long timeoutMs,
boolean failOpen,
+ LongSupplier currentTimeMs, AuthProvider authProvider) {
+ this(context._serverInstanceId, context._onlineBrokersSupplier,
context._allBrokersReady,
+ context._checkExecutor, context._routingStatusClient, timeoutMs,
failOpen, currentTimeMs,
+ authProvider, context._state);
+ }
+
+ @VisibleForTesting
+ BrokerRoutingReadyChecker(String serverInstanceId, Supplier<Set<String>>
onlineBrokersSupplier,
+ Predicate<Set<String>> allBrokersReady) {
+ this(serverInstanceId, onlineBrokersSupplier, allBrokersReady,
Long.MAX_VALUE, false, () -> 0L);
+ }
+
+ @VisibleForTesting
+ BrokerRoutingReadyChecker(String serverInstanceId, Supplier<Set<String>>
onlineBrokersSupplier,
+ Predicate<Set<String>> allBrokersReady, long timeoutMs, boolean
failOpen, LongSupplier currentTimeMs) {
+ this(serverInstanceId, onlineBrokersSupplier, (brokers, shouldStop) ->
allBrokersReady.test(brokers), null,
+ RoutingStatusClient.NOOP, timeoutMs, failOpen, currentTimeMs, new
NullAuthProvider(),
+ new AtomicReference<>(State.CHECKING));
+ }
+
+ private BrokerRoutingReadyChecker(String serverInstanceId,
Supplier<Set<String>> onlineBrokersSupplier,
+ BrokersReadyEvaluator allBrokersReady, @Nullable
ScheduledExecutorService checkExecutor,
+ RoutingStatusClient routingStatusClient, long timeoutMs, boolean
failOpen, LongSupplier currentTimeMs,
+ AuthProvider authProvider, AtomicReference<State> state) {
+ _serverInstanceId = serverInstanceId;
+ _onlineBrokersSupplier = onlineBrokersSupplier;
+ _allBrokersReady = allBrokersReady;
+ _checkExecutor = checkExecutor;
+ _routingStatusClient = routingStatusClient;
+ _currentTimeMs = currentTimeMs;
+ long nowMs = currentTimeMs.getAsLong();
+ _deadlineMs = timeoutMs >= Long.MAX_VALUE - nowMs ? Long.MAX_VALUE : nowMs
+ Math.max(timeoutMs, 0L);
+ _failOpen = failOpen;
+ _authProvider = authProvider;
+ _state = state;
+ if (_checkExecutor != null) {
+ _checkExecutor.scheduleWithFixedDelay(this::check, 0L,
CHECK_INTERVAL_MS, TimeUnit.MILLISECONDS);
+ }
+ }
+
+ public boolean isReady() {
+ transitionToReadyOnFailOpenTimeout();
+ return _state.get() == State.READY;
+ }
+
+ @VisibleForTesting
+ AuthProvider getAuthProvider() {
+ return _authProvider;
+ }
+
+ @VisibleForTesting
+ synchronized void check() {
+ if (_state.get() != State.CHECKING ||
transitionToReadyOnFailOpenTimeout()) {
+ return;
+ }
+ try {
+ Set<String> onlineBrokers = _onlineBrokersSupplier.get();
+ if (!onlineBrokers.isEmpty() && _allBrokersReady.test(onlineBrokers,
this::shouldStopChecking)
+ && _state.get() == State.CHECKING
+ // Do not mark the server ready if broker membership changed while
acknowledgements were collected.
+ && onlineBrokers.equals(_onlineBrokersSupplier.get())) {
+ if (_state.compareAndSet(State.CHECKING, State.READY)) {
+ LOGGER.info("All online brokers report server {} as routable: {}",
_serverInstanceId, onlineBrokers);
+ stopChecking();
+ }
+ return;
+ }
+ } catch (Exception e) {
+ LOGGER.debug("Failed to check broker routing readiness for server {}",
_serverInstanceId, e);
+ }
+
+ if (transitionToReadyOnFailOpenTimeout()) {
+ return;
+ }
+ if (_currentTimeMs.getAsLong() >= _deadlineMs) {
+ if (!_timeoutLogged) {
+ LOGGER.warn("Timed out waiting for all online brokers to report server
{} as routable; failOpen={}",
+ _serverInstanceId, _failOpen);
+ _timeoutLogged = true;
+ }
+ }
+ }
+
+ private boolean shouldStopChecking() {
+ return _state.get() != State.CHECKING ||
Thread.currentThread().isInterrupted()
+ || transitionToReadyOnFailOpenTimeout();
+ }
+
+ private boolean transitionToReadyOnFailOpenTimeout() {
+ if (!_failOpen || _currentTimeMs.getAsLong() < _deadlineMs) {
+ return false;
+ }
+ if (_state.compareAndSet(State.CHECKING, State.READY)) {
+ LOGGER.warn("Timed out waiting for all online brokers to report server
{} as routable; failOpen=true",
+ _serverInstanceId);
+ stopChecking();
+ }
+ return _state.get() != State.CHECKING;
+ }
+
+ private void stopChecking() {
+ if (_checkExecutor != null) {
+ _checkExecutor.shutdown();
+ }
+ }
+
+ private static ProductionContext createProductionContext(HelixManager
helixManager, AuthProvider authProvider) {
+ HttpClientConfig httpClientConfig = HttpClientConfig.newBuilder()
+ .withMaxConns(1)
+ .withMaxConnsPerRoute(1)
+ .withConnectionTimeoutMs((int) REQUEST_TIMEOUT_MS)
+ .withFollowRedirects(false)
+ .build();
+ RoutingStatusClient routingStatusClient = new SecureRoutingStatusClient(
+ new HttpRoutingStatusClient(new HttpClient(httpClientConfig,
TlsUtils.getSslContext(), true)));
+ ScheduledExecutorService checkExecutor =
Executors.newSingleThreadScheduledExecutor(
+ new
ThreadFactoryBuilder().setNameFormat("broker-routing-ready-check-%d").setDaemon(true).build());
+ return createContext(helixManager, routingStatusClient, checkExecutor,
authProvider,
+ new AtomicReference<>(State.CHECKING));
+ }
+
+ private static ProductionContext createContext(HelixManager helixManager,
RoutingStatusClient routingStatusClient,
+ @Nullable ScheduledExecutorService checkExecutor, AuthProvider
authProvider, AtomicReference<State> state) {
+ String serverInstanceId = helixManager.getInstanceName();
+ HelixAdmin helixAdmin = helixManager.getClusterManagmentTool();
+ String clusterName = helixManager.getClusterName();
+ Supplier<Set<String>> onlineBrokersSupplier = () -> {
+ ExternalView brokerResource =
helixAdmin.getResourceExternalView(clusterName,
+ CommonConstants.Helix.BROKER_RESOURCE_INSTANCE);
+ return HelixHelper.getOnlineInstanceFromExternalView(brokerResource);
+ };
+ RoutingStatusClient secureRoutingStatusClient =
+ routingStatusClient instanceof SecureRoutingStatusClient ?
routingStatusClient
+ : new SecureRoutingStatusClient(routingStatusClient);
+ BrokersReadyEvaluator allBrokersReady = (brokers, shouldStop) -> {
+ for (String broker : brokers) {
+ if (shouldStop.getAsBoolean() || state.get() != State.CHECKING ||
Thread.currentThread().isInterrupted()) {
Review Comment:
[Simplification] The shouldStop callback already checks _state, thread
interruption, and the fail-open deadline. The two checks beside it repeat its
first two conditions. Using only shouldStop here would also remove the need to
pass the state reference through createContext and ProductionContext.
##########
pinot-common/src/main/java/org/apache/pinot/common/utils/http/HttpClient.java:
##########
@@ -343,13 +398,58 @@ public SimpleHttpResponse sendMultipartPostRequest(String
url, String body, @Nul
private static String httpEntityToString(HttpEntity httpEntity)
throws IOException {
+ return httpEntityToString(httpEntity, Integer.MAX_VALUE);
+ }
+
+ private static String httpEntityToString(HttpEntity httpEntity, int
maxResponseLength)
+ throws IOException {
try {
- return EntityUtils.toString(httpEntity);
+ return EntityUtils.toString(httpEntity, maxResponseLength);
} catch (ParseException exception) {
throw new RuntimeException(exception);
}
}
+ private static BoundedResponseContent readResponseContent(HttpEntity
httpEntity, int maxResponseLength)
+ throws IOException {
+ Preconditions.checkArgument(maxResponseLength >= 0, "Maximum response
length must be non-negative");
+ if (maxResponseLength == Integer.MAX_VALUE) {
Review Comment:
[MAJOR] Default sendRequest calls pass Integer.MAX_VALUE, so this branch
converts the entity before the null check below. A 304 response has no entity:
this now throws IllegalArgumentException instead of returning the HTTP status,
which the previous error formatter did. Checking httpEntity == null first would
also cover bodyless HEAD error responses.
--
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]