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 67aa25cef0d fix: Avoid sleep on stop in K8sDruidNodeDiscoveryProvider.
(#19748)
67aa25cef0d is described below
commit 67aa25cef0d0788b89062838871e2964e3a52285
Author: Gian Merlino <[email protected]>
AuthorDate: Sat Jul 25 04:11:04 2026 -0700
fix: Avoid sleep on stop in K8sDruidNodeDiscoveryProvider. (#19748)
Prior to this patch, an InterruptedException caused by shutdownNow()
in stop() would trigger sleep(watcherErrorRetryWaitMS), leading to a
ten second delay on stop().
---
.../discovery/K8sDruidNodeDiscoveryProvider.java | 18 +++--
.../K8sDruidNodeDiscoveryProviderTest.java | 94 ++++++++++++++++++++++
2 files changed, 106 insertions(+), 6 deletions(-)
diff --git
a/extensions-core/kubernetes-extensions/src/main/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProvider.java
b/extensions-core/kubernetes-extensions/src/main/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProvider.java
index 96497d4510f..8dfc595db74 100644
---
a/extensions-core/kubernetes-extensions/src/main/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProvider.java
+++
b/extensions-core/kubernetes-extensions/src/main/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProvider.java
@@ -223,7 +223,7 @@ public class K8sDruidNodeDiscoveryProvider extends
DruidNodeDiscoveryProvider
return;
}
- while (lifecycleLock.awaitStarted(1, TimeUnit.MILLISECONDS)) {
+ while (lifecycleLock.isStarted()) {
try {
DiscoveryDruidNodeList list =
k8sApiClient.listPods(podInfo.getPodNamespace(), labelSelector, nodeRole);
baseNodeRoleWatcher.resetNodes(list.getDruidNodes());
@@ -241,8 +241,10 @@ public class K8sDruidNodeDiscoveryProvider extends
DruidNodeDiscoveryProvider
catch (Throwable ex) {
LOGGER.error(ex, "Exception while watching for role[%s].", nodeRole);
- // Wait a little before trying again.
- sleep(watcherErrorRetryWaitMS);
+ if (lifecycleLock.isStarted()) {
+ // If not stopped, wait a little before trying again.
+ sleep(watcherErrorRetryWaitMS);
+ }
}
}
@@ -252,7 +254,7 @@ public class K8sDruidNodeDiscoveryProvider extends
DruidNodeDiscoveryProvider
private void keepWatching(String labelSelector, String resourceVersion)
{
String nextResourceVersion = resourceVersion;
- while (lifecycleLock.awaitStarted(1, TimeUnit.MILLISECONDS)) {
+ while (lifecycleLock.isStarted()) {
try {
WatchResult iter =
k8sApiClient.watchPods(podInfo.getPodNamespace(), labelSelector,
nextResourceVersion, nodeRole);
@@ -300,11 +302,15 @@ public class K8sDruidNodeDiscoveryProvider extends
DruidNodeDiscoveryProvider
catch (SocketTimeoutException ex) {
// socket read timeout can happen normally due to k8s not having
anything new to push leading to socket
// read timeout, so no error log
- sleep(watcherErrorRetryWaitMS);
+ if (lifecycleLock.isStarted()) {
+ sleep(watcherErrorRetryWaitMS);
+ }
}
catch (Throwable ex) {
LOGGER.error(ex, "Error while watching role[%s]", this.nodeRole);
- sleep(watcherErrorRetryWaitMS);
+ if (lifecycleLock.isStarted()) {
+ sleep(watcherErrorRetryWaitMS);
+ }
}
}
}
diff --git
a/extensions-core/kubernetes-extensions/src/test/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProviderTest.java
b/extensions-core/kubernetes-extensions/src/test/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProviderTest.java
index 18e030f2507..6417edb4b27 100644
---
a/extensions-core/kubernetes-extensions/src/test/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProviderTest.java
+++
b/extensions-core/kubernetes-extensions/src/test/java/org/apache/druid/k8s/discovery/K8sDruidNodeDiscoveryProviderTest.java
@@ -40,6 +40,7 @@ import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
public class K8sDruidNodeDiscoveryProviderTest
@@ -430,6 +431,61 @@ public class K8sDruidNodeDiscoveryProviderTest
discoveryProvider.stop();
}
+ @Test
+ @Timeout(value = 60_000, unit = TimeUnit.MILLISECONDS)
+ public void testStopDoesNotWaitForErrorRetryWhenWatchIsInterrupted() throws
Exception
+ {
+ final String labelSelector =
"druidDiscoveryAnnouncement-cluster-identifier=druid-cluster,druidDiscoveryAnnouncement-router=true";
+
+ // Long enough that a single unguarded retry sleep outlasts the 15s stop()
timeout on watchExecutor.
+ final long watcherErrorRetryWaitMS = 30_000;
+
+ final CountDownLatch watchBlocked = new CountDownLatch(1);
+
+ final K8sApiClient mockK8sApiClient =
EasyMock.createMock(K8sApiClient.class);
+ EasyMock.expect(mockK8sApiClient.listPods(podInfo.getPodNamespace(),
labelSelector, NodeRole.ROUTER)).andReturn(
+ new DiscoveryDruidNodeList(
+ "v1",
+ ImmutableMap.of(testNode1.getDruidNode().getHostAndPortToUse(),
testNode1)
+ )
+ );
+ EasyMock.expect(mockK8sApiClient.watchPods(
+ podInfo.getPodNamespace(), labelSelector, "v1",
NodeRole.ROUTER)).andReturn(
+ new BlockingWatchResult(watchBlocked)
+ );
+ EasyMock.replay(mockK8sApiClient);
+
+ final K8sDruidNodeDiscoveryProvider discoveryProvider = new
K8sDruidNodeDiscoveryProvider(
+ podInfo,
+ discoveryConfig,
+ mockK8sApiClient,
+ watcherErrorRetryWaitMS
+ );
+ discoveryProvider.start();
+
+ final K8sDruidNodeDiscoveryProvider.NodeRoleWatcher nodeDiscovery =
+ discoveryProvider.getForNodeRole(NodeRole.ROUTER, false);
+ nodeDiscovery.start();
+
+ Assertions.assertTrue(
+ watchBlocked.await(30, TimeUnit.SECONDS),
+ "watch thread never reached hasNext()"
+ );
+
+ final long startNs = System.nanoTime();
+ discoveryProvider.stop();
+ final long stopMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() -
startNs);
+
+ Assertions.assertTrue(
+ stopMillis < 5_000,
+ StringUtils.format(
+ "stop() took [%,d] ms, expected it not to wait out
watcherErrorRetryWaitMS[%,d]",
+ stopMillis,
+ watcherErrorRetryWaitMS
+ )
+ );
+ }
+
private static class MockListener implements DruidNodeDiscovery.Listener
{
List<Event> events;
@@ -554,6 +610,44 @@ public class K8sDruidNodeDiscoveryProviderTest
}
}
+ /**
+ * A {@link WatchResult} that parks in {@link #hasNext()} the way an open
k8s watch socket does, and surfaces the
+ * shutdown interrupt as a {@link RuntimeException}.
+ */
+ private static class BlockingWatchResult implements WatchResult
+ {
+ private final CountDownLatch blocked;
+
+ public BlockingWatchResult(CountDownLatch blocked)
+ {
+ this.blocked = blocked;
+ }
+
+ @Override
+ public boolean hasNext()
+ {
+ blocked.countDown();
+ try {
+ Thread.sleep(Long.MAX_VALUE);
+ return false; // Not reached.
+ }
+ catch (InterruptedException ex) {
+ throw new RuntimeException(ex);
+ }
+ }
+
+ @Override
+ public Watch.Response<DiscoveryDruidNodeAndResourceVersion> next()
+ {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public void close()
+ {
+ }
+ }
+
private static class MockWatchResult implements WatchResult
{
private List<Watch.Response<DiscoveryDruidNodeAndResourceVersion>> results;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]