This is an automated email from the ASF dual-hosted git repository.

cryptoe pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new a92f4b62b8d fix: Don't reopen closed clients in DartWorkerClientImpl. 
(#19927)
a92f4b62b8d is described below

commit a92f4b62b8d8241c1d3f7a9259c1e69697dea6da
Author: Gian Merlino <[email protected]>
AuthorDate: Tue Aug 11 07:14:04 2026 -0700

    fix: Don't reopen closed clients in DartWorkerClientImpl. (#19927)
    
    DartWorkerClientImpl is scoped to a single query, so there is no need
    to be able to reopen clients (once a worker fails, the query also
    fails).
    
    This patch fixes a retry loop that could be caused when a server goes
    away: ControllerMessageListener#serverRemoved calls closeClient, but
    then the next call to the worker would cause the client to be re-created
    and keep retrying until its retries are exhausted.
---
 .../dart/controller/ControllerMessageListener.java |   2 +-
 .../msq/dart/controller/DartWorkerManager.java     |  13 ++-
 .../druid/msq/dart/worker/DartWorkerClient.java    |   4 +-
 .../msq/dart/worker/DartWorkerClientImpl.java      |  49 ++++----
 .../msq/dart/worker/DartWorkerClientImplTest.java  | 124 +++++++++++++++++++++
 5 files changed, 165 insertions(+), 27 deletions(-)

diff --git 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/ControllerMessageListener.java
 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/ControllerMessageListener.java
index 0fe85f42c75..b6509269173 100644
--- 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/ControllerMessageListener.java
+++ 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/ControllerMessageListener.java
@@ -70,7 +70,7 @@ public class ControllerMessageListener implements 
MessageListener<ControllerMess
       // Close the worker client for this server.
       final ControllerContext controllerContext = 
controller.getControllerContext();
       if (controllerContext instanceof DartControllerContext) {
-        ((DartControllerContext) 
controllerContext).newWorkerClient().closeClient(workerId.getHostAndPort());
+        ((DartControllerContext) 
controllerContext).newWorkerClient().closeClient(workerIdString);
       }
 
       // Notify the controller that the worker has gone offline.
diff --git 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/DartWorkerManager.java
 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/DartWorkerManager.java
index f6d68d57f42..5df0f347f59 100644
--- 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/DartWorkerManager.java
+++ 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/controller/DartWorkerManager.java
@@ -31,7 +31,6 @@ import org.apache.druid.common.guava.FutureUtils;
 import org.apache.druid.error.DruidException;
 import org.apache.druid.indexer.TaskState;
 import org.apache.druid.java.util.common.ISE;
-import org.apache.druid.java.util.common.logger.Logger;
 import org.apache.druid.msq.dart.worker.DartWorkerClient;
 import org.apache.druid.msq.exec.ControllerContext;
 import org.apache.druid.msq.exec.WorkerClient;
@@ -39,7 +38,6 @@ import org.apache.druid.msq.exec.WorkerManager;
 import org.apache.druid.msq.exec.WorkerStats;
 import org.apache.druid.msq.indexing.WorkerCount;
 import org.apache.druid.msq.indexing.error.MSQFault;
-import org.apache.druid.utils.CloseableUtils;
 
 import java.util.ArrayList;
 import java.util.Collections;
@@ -56,8 +54,6 @@ import java.util.concurrent.atomic.AtomicReference;
  */
 public class DartWorkerManager implements WorkerManager
 {
-  private static final Logger log = new Logger(DartWorkerManager.class);
-
   private final List<String> workerIds;
   private final List<String> workerDescs;
   private final DartWorkerClient workerClient;
@@ -72,6 +68,14 @@ public class DartWorkerManager implements WorkerManager
     STOPPED
   }
 
+  /**
+   * Creates a new worker manager.
+   *
+   * @param workerIds    Fixed list of IDs of the workers to manage.
+   * @param workerDescs  Descriptions of the workers, same length as {@code 
workerIds}
+   * @param workerClient Client to use to contact workers. Not owned by this 
class. It should be closed externally
+   *                     after you are done using this class.
+   */
   public DartWorkerManager(
       final List<String> workerIds,
       final List<String> workerDescs,
@@ -210,7 +214,6 @@ public class DartWorkerManager implements WorkerManager
         }
       }
 
-      CloseableUtils.closeAndSuppressExceptions(workerClient, e -> log.warn(e, 
"Failed to close workerClient"));
       stopFuture.set(null);
     }
   }
diff --git 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClient.java
 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClient.java
index e048f9fd6d6..eb918eeed5d 100644
--- 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClient.java
+++ 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClient.java
@@ -28,9 +28,9 @@ public interface DartWorkerClient extends WorkerClient
   /**
    * Close a single worker's clients. Used when that worker fails, so we stop 
trying to contact it.
    *
-   * @param workerHost worker host:port
+   * @param workerId worker ID string
    */
-  void closeClient(String hostAndPort);
+  void closeClient(String workerId);
 
   /**
    * Stops a worker. Dart-only API, used by the {@link DartWorkerManager}.
diff --git 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClientImpl.java
 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClientImpl.java
index e1de222e1e0..8b73d768227 100644
--- 
a/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClientImpl.java
+++ 
b/multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClientImpl.java
@@ -30,6 +30,7 @@ import 
org.apache.druid.java.util.http.client.response.HttpResponseHandler;
 import org.apache.druid.msq.dart.worker.http.DartWorkerResource;
 import org.apache.druid.msq.exec.WorkerClient;
 import org.apache.druid.msq.rpc.BaseWorkerClientImpl;
+import org.apache.druid.query.QueryContexts;
 import org.apache.druid.rpc.FixedServiceLocator;
 import org.apache.druid.rpc.IgnoreHttpResponseHandler;
 import org.apache.druid.rpc.RequestBuilder;
@@ -64,10 +65,13 @@ public class DartWorkerClientImpl extends 
BaseWorkerClientImpl implements DartWo
   @GuardedBy("clientMap")
   private final Map<String, Pair<ServiceClient, Closeable>> clientMap = new 
HashMap<>();
 
+  @GuardedBy("clientMap")
+  private boolean closed;
+
   /**
    * Create a worker client.
    *
-   * @param queryId        dart query ID. see {@link 
org.apache.druid.query.QueryContexts#CTX_DART_QUERY_ID}
+   * @param queryId        dart query ID. see {@link 
QueryContexts#CTX_DART_QUERY_ID}
    * @param clientFactory  service client factor
    * @param smileMapper    Smile object mapper
    * @param controllerHost Controller host (see {@link 
DartWorkerResource#HEADER_CONTROLLER_HOST}) if this is a
@@ -99,26 +103,19 @@ public class DartWorkerClientImpl extends 
BaseWorkerClientImpl implements DartWo
   @Override
   protected ServiceClient getClient(final String workerIdString)
   {
-    final WorkerId workerId = WorkerId.fromString(workerIdString);
-    if (!queryId.equals(workerId.getQueryId())) {
-      throw DruidException.defensive("Unexpected queryId[%s]. Expected 
queryId[%s]", workerId.getQueryId(), queryId);
-    }
-
-    synchronized (clientMap) {
-      return clientMap.computeIfAbsent(workerId.getHostAndPort(), ignored -> 
makeNewClient(workerId)).left();
-    }
+    return getClientAndLocator(workerIdString).left();
   }
 
-  /**
-   * {@inheritDoc}
-   */
   @Override
-  public void closeClient(final String workerHost)
+  public void closeClient(final String workerIdString)
   {
+    // Close, but do not remove from clientMap, so it stays closed. Note that 
if closeClient is called before
+    // getClient(), then the client will be created and immediately closed. 
This is intentional, since it allows
+    // server-removed notifications to be respected even if we haven't tried 
to contact a worker yet.
     synchronized (clientMap) {
-      final Pair<ServiceClient, Closeable> clientPair = 
clientMap.remove(workerHost);
-      if (clientPair != null) {
-        CloseableUtils.closeAndWrapExceptions(clientPair.right());
+      // Do nothing if we have already been closed; in this case we know there 
are no clients active.
+      if (!closed) {
+        
CloseableUtils.closeAndWrapExceptions(getClientAndLocator(workerIdString).right());
       }
     }
   }
@@ -138,12 +135,10 @@ public class DartWorkerClientImpl extends 
BaseWorkerClientImpl implements DartWo
       }
 
       clientMap.clear();
+      closed = true;
     }
   }
 
-  /**
-   * {@inheritDoc}
-   */
   @Override
   public ListenableFuture<?> stopWorker(String workerId)
   {
@@ -173,6 +168,22 @@ public class DartWorkerClientImpl extends 
BaseWorkerClientImpl implements DartWo
     return Pair.of(client, locator);
   }
 
+  private Pair<ServiceClient, Closeable> getClientAndLocator(final String 
workerIdString)
+  {
+    final WorkerId workerId = WorkerId.fromString(workerIdString);
+    if (!queryId.equals(workerId.getQueryId())) {
+      throw DruidException.defensive("Unexpected queryId[%s]. Expected 
queryId[%s]", workerId.getQueryId(), queryId);
+    }
+
+    synchronized (clientMap) {
+      if (closed) {
+        throw DruidException.defensive("%s is closed", getClass().getName());
+      }
+
+      return clientMap.computeIfAbsent(workerId.getHostAndPort(), ignored -> 
makeNewClient(workerId));
+    }
+  }
+
   /**
    * Service client that adds the {@link 
DartWorkerResource#HEADER_CONTROLLER_HOST} header.
    */
diff --git 
a/multi-stage-query/src/test/java/org/apache/druid/msq/dart/worker/DartWorkerClientImplTest.java
 
b/multi-stage-query/src/test/java/org/apache/druid/msq/dart/worker/DartWorkerClientImplTest.java
new file mode 100644
index 00000000000..1c0a39aba2f
--- /dev/null
+++ 
b/multi-stage-query/src/test/java/org/apache/druid/msq/dart/worker/DartWorkerClientImplTest.java
@@ -0,0 +1,124 @@
+/*
+ * 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.druid.msq.dart.worker;
+
+import org.apache.druid.error.DruidException;
+import org.apache.druid.java.util.common.concurrent.Execs;
+import org.apache.druid.java.util.http.client.HttpClient;
+import org.apache.druid.rpc.ServiceClient;
+import org.apache.druid.rpc.ServiceClientFactoryImpl;
+import org.apache.druid.rpc.ServiceClosedException;
+import org.apache.druid.segment.TestHelper;
+import org.hamcrest.CoreMatchers;
+import org.hamcrest.MatcherAssert;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ScheduledExecutorService;
+
+public class DartWorkerClientImplTest
+{
+  private static final String QUERY_ID = "abc-123";
+  private static final WorkerId WORKER_ID = new WorkerId("http", 
"localhost:8100", QUERY_ID);
+
+  private ScheduledExecutorService connectExec;
+  private DartWorkerClientImpl workerClient;
+
+  @BeforeEach
+  public void setUp()
+  {
+    connectExec = Execs.scheduledSingleThreaded("DartWorkerClientImplTest-%s");
+    workerClient = new DartWorkerClientImpl(
+        QUERY_ID,
+        new ServiceClientFactoryImpl(Mockito.mock(HttpClient.class), 
connectExec),
+        TestHelper.makeSmileMapper(),
+        "localhost:8080"
+    );
+  }
+
+  @AfterEach
+  public void tearDown()
+  {
+    workerClient.close();
+    connectExec.shutdownNow();
+  }
+
+  @Test
+  public void test_getClient_isCachedPerWorker()
+  {
+    final ServiceClient client = workerClient.getClient(WORKER_ID.toString());
+    Assertions.assertSame(client, 
workerClient.getClient(WORKER_ID.toString()));
+  }
+
+  @Test
+  public void test_getClient_wrongQueryId()
+  {
+    final WorkerId otherWorkerId = new WorkerId("http", "localhost:8100", 
"other-query");
+    Assertions.assertThrows(
+        DruidException.class,
+        () -> workerClient.getClient(otherWorkerId.toString())
+    );
+  }
+
+  @Test
+  public void test_closeClient_staysClosed()
+  {
+    final ServiceClient client = workerClient.getClient(WORKER_ID.toString());
+    workerClient.closeClient(WORKER_ID.toString());
+
+    // The closed client is retained, rather than being replaced by a fresh 
one that would contact the worker again.
+    Assertions.assertSame(client, 
workerClient.getClient(WORKER_ID.toString()));
+    assertRequestFailsAsClosed();
+  }
+
+  @Test
+  public void test_closeClient_beforeGetClient()
+  {
+    // Closing a worker we never contacted still prevents it from being 
contacted later.
+    workerClient.closeClient(WORKER_ID.toString());
+    assertRequestFailsAsClosed();
+  }
+
+  @Test
+  public void test_closeClient_afterClose_isNoop()
+  {
+    workerClient.close();
+    Assertions.assertDoesNotThrow(() -> 
workerClient.closeClient(WORKER_ID.toString()));
+    Assertions.assertThrows(DruidException.class, () -> 
workerClient.getClient(WORKER_ID.toString()));
+  }
+
+  /**
+   * Verify that a request to {@link #WORKER_ID} fails immediately, rather 
than retrying, due to its client
+   * being closed.
+   */
+  private void assertRequestFailsAsClosed()
+  {
+    final ExecutionException e = Assertions.assertThrows(
+        ExecutionException.class,
+        () -> workerClient.stopWorker(WORKER_ID.toString()).get()
+    );
+
+    MatcherAssert.assertThat(e.getCause(), 
CoreMatchers.instanceOf(ServiceClosedException.class));
+  }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to