From 53a5e09f08b8c81a7d1d3370a08e922a38a2ddd1 Mon Sep 17 00:00:00 2001 From: Shrey Narayan Date: Thu, 20 Aug 2026 02:21:17 -0700 Subject: [PATCH 1/2] SOLR-18298: only run ZkController reconnect recovery after ZooKeeper session expiry. Curator RECONNECTED fires on every ZK instance hop, which made rolling ZK restarts re-elect leaders and re-register cores. Restore Solr 9 behavior by treating LOST then RECONNECTED as expiration, authored by Shrey Narayan (NextBrick). Co-authored-by: Cursor --- ...SOLR-18298-zk-reconnect-session-expiry.yml | 8 +++ .../org/apache/solr/cloud/ZkController.java | 3 + .../solr/common/cloud/OnDisconnect.java | 12 +++- .../apache/solr/common/cloud/OnReconnect.java | 27 ++++++- .../cloud/TestOnReconnectSessionExpiry.java | 70 +++++++++++++++++++ 5 files changed, 117 insertions(+), 3 deletions(-) create mode 100644 changelog/unreleased/SOLR-18298-zk-reconnect-session-expiry.yml create mode 100644 solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java diff --git a/changelog/unreleased/SOLR-18298-zk-reconnect-session-expiry.yml b/changelog/unreleased/SOLR-18298-zk-reconnect-session-expiry.yml new file mode 100644 index 000000000000..e9f6f4527c43 --- /dev/null +++ b/changelog/unreleased/SOLR-18298-zk-reconnect-session-expiry.yml @@ -0,0 +1,8 @@ +title: ZkController no longer treats every ZooKeeper reconnect as session expiration. Re-election and core re-registration now run only after Curator reports ConnectionState.LOST, restoring Solr 9 behavior during ZK rolling restarts. +type: fixed +authors: + - name: Shrey Narayan (NextBrick) + url: https://nextbrick.com +links: + - name: SOLR-18298 + url: https://issues.apache.org/jira/browse/SOLR-18298 diff --git a/solr/core/src/java/org/apache/solr/cloud/ZkController.java b/solr/core/src/java/org/apache/solr/cloud/ZkController.java index 3ab82e68b1d6..c94218056294 100644 --- a/solr/core/src/java/org/apache/solr/cloud/ZkController.java +++ b/solr/core/src/java/org/apache/solr/cloud/ZkController.java @@ -402,6 +402,9 @@ public ZkController( } private void onDisconnect(boolean sessionExpired) { + if (!sessionExpired) { + return; + } try { overseer.close(); } catch (Exception e) { diff --git a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java index 9535a59cef55..b178c41c9aed 100644 --- a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java +++ b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java @@ -20,13 +20,21 @@ import org.apache.curator.framework.state.ConnectionState; import org.apache.curator.framework.state.ConnectionStateListener; +/** + * Listener for ZooKeeper session loss. + * + *

When registered as a Curator {@link ConnectionStateListener}, {@link #onDisconnect(boolean)} + * runs only after session expiration ({@link ConnectionState#LOST}). A transient {@link + * ConnectionState#SUSPENDED} keeps the session and must not tear down SolrCloud leadership. See + * SOLR-18298. + */ public interface OnDisconnect extends ConnectionStateListener { void onDisconnect(boolean sessionExpired); @Override default void stateChanged(CuratorFramework client, ConnectionState newState) { - if (newState == ConnectionState.LOST || newState == ConnectionState.SUSPENDED) { - onDisconnect(newState == ConnectionState.LOST); + if (newState == ConnectionState.LOST) { + onDisconnect(true); } } } diff --git a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java index 8d54312d3e0f..88f95efffea7 100644 --- a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java +++ b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java @@ -16,6 +16,9 @@ */ package org.apache.solr.common.cloud; +import java.util.Collections; +import java.util.Set; +import java.util.WeakHashMap; import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.state.ConnectionState; import org.apache.curator.framework.state.ConnectionStateListener; @@ -26,14 +29,36 @@ * implementation should call * org.apache.solr.cloud.ZkController#removeOnReconnectListener(OnReconnect) when it no longer needs * to be notified of ZK reconnection events. + * + *

When registered as a Curator {@link ConnectionStateListener}, {@link #onReconnect()} runs only + * after a session expiration ({@link ConnectionState#LOST} then {@link + * ConnectionState#RECONNECTED}). A reconnect after {@link ConnectionState#SUSPENDED} keeps the + * ZooKeeper session and must not trigger SolrCloud recovery. See SOLR-18298. */ public interface OnReconnect extends ConnectionStateListener { void onReconnect(); @Override default void stateChanged(CuratorFramework client, ConnectionState newState) { - if (ConnectionState.RECONNECTED.equals(newState)) { + if (newState == ConnectionState.LOST) { + LostSessions.mark(this); + } else if (newState == ConnectionState.RECONNECTED && LostSessions.consume(this)) { onReconnect(); } } + + /** Tracks listeners that have observed {@link ConnectionState#LOST} and still need reconnect. */ + final class LostSessions { + private static final Set lost = Collections.newSetFromMap(new WeakHashMap<>()); + + private LostSessions() {} + + static synchronized void mark(OnReconnect listener) { + lost.add(listener); + } + + static synchronized boolean consume(OnReconnect listener) { + return lost.remove(listener); + } + } } diff --git a/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java b/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java new file mode 100644 index 000000000000..9d6c83fa65af --- /dev/null +++ b/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java @@ -0,0 +1,70 @@ +/* + * 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.solr.common.cloud; + +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.curator.framework.state.ConnectionState; +import org.apache.solr.SolrTestCase; +import org.junit.Test; + +/** SOLR-18298: Curator RECONNECTED is not the same as Solr 9 session-expiration reconnect. */ +public class TestOnReconnectSessionExpiry extends SolrTestCase { + + @Test + public void testReconnectAfterSuspendDoesNotFire() { + AtomicInteger reconnects = new AtomicInteger(); + OnReconnect listener = () -> reconnects.incrementAndGet(); + + listener.stateChanged(null, ConnectionState.SUSPENDED); + listener.stateChanged(null, ConnectionState.RECONNECTED); + + assertEquals(0, reconnects.get()); + } + + @Test + public void testReconnectAfterLostDoesFire() { + AtomicInteger reconnects = new AtomicInteger(); + OnReconnect listener = () -> reconnects.incrementAndGet(); + + listener.stateChanged(null, ConnectionState.LOST); + listener.stateChanged(null, ConnectionState.RECONNECTED); + + assertEquals(1, reconnects.get()); + } + + @Test + public void testReconnectWithoutPriorLostDoesNotFire() { + AtomicInteger reconnects = new AtomicInteger(); + OnReconnect listener = () -> reconnects.incrementAndGet(); + + listener.stateChanged(null, ConnectionState.RECONNECTED); + + assertEquals(0, reconnects.get()); + } + + @Test + public void testDisconnectFiresOnlyOnLost() { + AtomicInteger disconnects = new AtomicInteger(); + OnDisconnect listener = sessionExpired -> disconnects.incrementAndGet(); + + listener.stateChanged(null, ConnectionState.SUSPENDED); + assertEquals(0, disconnects.get()); + + listener.stateChanged(null, ConnectionState.LOST); + assertEquals(1, disconnects.get()); + } +} From 21cb72bc750de09961f5f8efee711a5ad73a9cc0 Mon Sep 17 00:00:00 2001 From: rayshrey <121871912+rayshrey@users.noreply.github.com> Date: Wed, 26 Aug 2026 15:56:46 -0700 Subject: [PATCH 2/2] SOLR-18298: Track session expiry in owning components --- .../org/apache/solr/cloud/ZkController.java | 6 +++ .../apache/solr/cloud/ZkControllerTest.java | 51 +++++++++++++++++++ .../solr/common/cloud/OnDisconnect.java | 12 +---- .../apache/solr/common/cloud/OnReconnect.java | 27 +--------- .../solr/common/cloud/ZkStateReader.java | 42 ++++++++------- .../cloud/TestOnReconnectSessionExpiry.java | 44 +++++++--------- 6 files changed, 102 insertions(+), 80 deletions(-) diff --git a/solr/core/src/java/org/apache/solr/cloud/ZkController.java b/solr/core/src/java/org/apache/solr/cloud/ZkController.java index c94218056294..26d83e8a817c 100644 --- a/solr/core/src/java/org/apache/solr/cloud/ZkController.java +++ b/solr/core/src/java/org/apache/solr/cloud/ZkController.java @@ -54,6 +54,7 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Predicate; import java.util.stream.Collectors; @@ -214,6 +215,7 @@ public String toString() { new SolrNamedThreadFactory("zkConnectionListenerCallback")); private final OnReconnect onReconnect = this::onReconnect; private final OnDisconnect onDisconnect = this::onDisconnect; + private final AtomicBoolean zkSessionExpired = new AtomicBoolean(); private final String zkServerAddress; // example: 127.0.0.1:54062/solr @@ -405,6 +407,7 @@ private void onDisconnect(boolean sessionExpired) { if (!sessionExpired) { return; } + zkSessionExpired.set(true); try { overseer.close(); } catch (Exception e) { @@ -434,6 +437,9 @@ private T loadPluginOrDefault( } private void onReconnect() { + if (!zkSessionExpired.compareAndSet(true, false)) { + return; + } // on reconnect, reload cloud info log.info("ZooKeeper session re-connected ... refreshing core states after session expiration."); clearZkCollectionTerms(); diff --git a/solr/core/src/test/org/apache/solr/cloud/ZkControllerTest.java b/solr/core/src/test/org/apache/solr/cloud/ZkControllerTest.java index 5c7cd7442d3f..b33937c60418 100644 --- a/solr/core/src/test/org/apache/solr/cloud/ZkControllerTest.java +++ b/solr/core/src/test/org/apache/solr/cloud/ZkControllerTest.java @@ -36,6 +36,9 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; +import org.apache.curator.CuratorZookeeperClient; +import org.apache.curator.test.InstanceSpec; +import org.apache.curator.test.TestingCluster; import org.apache.solr.SolrTestCaseJ4; import org.apache.solr.client.api.util.SolrVersion; import org.apache.solr.client.solrj.jetty.HttpJettySolrClient; @@ -50,6 +53,7 @@ import org.apache.solr.common.cloud.ZkStateReader; import org.apache.solr.common.params.CollectionParams; import org.apache.solr.common.util.ExecutorUtil; +import org.apache.solr.common.util.RetryUtil; import org.apache.solr.common.util.SolrNamedThreadFactory; import org.apache.solr.common.util.Utils; import org.apache.solr.core.CloudConfig; @@ -766,6 +770,53 @@ public void testOverseerEnabledClusterPropertyTrue() throws Exception { } } + @Test + public void testReconnectRecoveryRequiresSessionExpiration() throws Exception { + try (TestingCluster zkCluster = new TestingCluster(3)) { + zkCluster.start(); + CoreContainer cc = getCoreContainer(); + try { + CloudConfig cloudConfig = new CloudConfig.CloudConfigBuilder("127.0.0.1", 8983).build(); + try (ZkController zkController = + new ZkController(cc, zkCluster.getConnectString(), TIMEOUT, cloudConfig)) { + AtomicInteger recoveries = new AtomicInteger(); + zkController.addOnReconnectListener(recoveries::incrementAndGet); + CuratorZookeeperClient curatorClient = + zkController.getZkClient().getCuratorFramework().getZookeeperClient(); + + InstanceSpec connected = zkCluster.findConnectionInstance(curatorClient.getZooKeeper()); + assertNotNull(connected); + zkCluster.killServer(connected); + RetryUtil.retryUntil( + "Solr did not connect to another ZooKeeper server", + 30, + 200, + TimeUnit.MILLISECONDS, + () -> zkCluster.findConnectionInstance(curatorClient.getZooKeeper()), + current -> current != null && !current.equals(connected)); + assertEquals( + "A transient ZooKeeper reconnect must not trigger session-expiration recovery", + 0, + recoveries.get()); + + curatorClient.getZooKeeper().getTestable().injectSessionExpiration(); + RetryUtil.retryUntil( + "A reconnect after session expiration did not trigger recovery", + 30, + 200, + TimeUnit.MILLISECONDS, + () -> recoveries.get() == 1); + assertEquals(1, recoveries.get()); + } + } finally { + cc.shutdown(); + } + } finally { + // TestingCluster closes its quorum asynchronously; allow its worker threads to terminate. + Thread.sleep(3000); + } + } + private CoreContainer getCoreContainer() { return new MockCoreContainer(); } diff --git a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java index b178c41c9aed..9535a59cef55 100644 --- a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java +++ b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnDisconnect.java @@ -20,21 +20,13 @@ import org.apache.curator.framework.state.ConnectionState; import org.apache.curator.framework.state.ConnectionStateListener; -/** - * Listener for ZooKeeper session loss. - * - *

When registered as a Curator {@link ConnectionStateListener}, {@link #onDisconnect(boolean)} - * runs only after session expiration ({@link ConnectionState#LOST}). A transient {@link - * ConnectionState#SUSPENDED} keeps the session and must not tear down SolrCloud leadership. See - * SOLR-18298. - */ public interface OnDisconnect extends ConnectionStateListener { void onDisconnect(boolean sessionExpired); @Override default void stateChanged(CuratorFramework client, ConnectionState newState) { - if (newState == ConnectionState.LOST) { - onDisconnect(true); + if (newState == ConnectionState.LOST || newState == ConnectionState.SUSPENDED) { + onDisconnect(newState == ConnectionState.LOST); } } } diff --git a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java index 88f95efffea7..8d54312d3e0f 100644 --- a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java +++ b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/OnReconnect.java @@ -16,9 +16,6 @@ */ package org.apache.solr.common.cloud; -import java.util.Collections; -import java.util.Set; -import java.util.WeakHashMap; import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.state.ConnectionState; import org.apache.curator.framework.state.ConnectionStateListener; @@ -29,36 +26,14 @@ * implementation should call * org.apache.solr.cloud.ZkController#removeOnReconnectListener(OnReconnect) when it no longer needs * to be notified of ZK reconnection events. - * - *

When registered as a Curator {@link ConnectionStateListener}, {@link #onReconnect()} runs only - * after a session expiration ({@link ConnectionState#LOST} then {@link - * ConnectionState#RECONNECTED}). A reconnect after {@link ConnectionState#SUSPENDED} keeps the - * ZooKeeper session and must not trigger SolrCloud recovery. See SOLR-18298. */ public interface OnReconnect extends ConnectionStateListener { void onReconnect(); @Override default void stateChanged(CuratorFramework client, ConnectionState newState) { - if (newState == ConnectionState.LOST) { - LostSessions.mark(this); - } else if (newState == ConnectionState.RECONNECTED && LostSessions.consume(this)) { + if (ConnectionState.RECONNECTED.equals(newState)) { onReconnect(); } } - - /** Tracks listeners that have observed {@link ConnectionState#LOST} and still need reconnect. */ - final class LostSessions { - private static final Set lost = Collections.newSetFromMap(new WeakHashMap<>()); - - private LostSessions() {} - - static synchronized void mark(OnReconnect listener) { - lost.add(listener); - } - - static synchronized boolean consume(OnReconnect listener) { - return lost.remove(listener); - } - } } diff --git a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/ZkStateReader.java b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/ZkStateReader.java index 4f0c3bb38366..01223a88509c 100644 --- a/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/ZkStateReader.java +++ b/solr/solrj-zookeeper/src/java/org/apache/solr/common/cloud/ZkStateReader.java @@ -390,6 +390,9 @@ private StatefulCollectionWatch(StateWatcher associatedWatcher) { private final SolrZkClient zkClient; private final boolean closeClient; + private final AtomicBoolean zkSessionExpired = new AtomicBoolean(); + private final OnDisconnect onDisconnect = this::onDisconnect; + private final OnReconnect onReconnect = this::onReconnect; private volatile boolean closed = false; @@ -424,29 +427,34 @@ public ZkStateReader( .withConnTimeOut(zkClientConnectTimeout, TimeUnit.MILLISECONDS) .withUseDefaultCredsAndACLs(canUseZkACLs) .build(); - this.zkClient - .getCuratorFramework() - .getConnectionStateListenable() - .addListener( - (OnReconnect) - () -> { - // on reconnect, reload cloud info - try { - this.createClusterStateWatchersAndUpdate(); - } catch (InterruptedException e) { - // Restore the interrupted status - Thread.currentThread().interrupt(); - log.warn("Interrupted", e); - } catch (Throwable e) { - log.error("An error has occurred while updating the cluster state", e); - } - }); + this.zkClient.getCuratorFramework().getConnectionStateListenable().addListener(onReconnect); + this.zkClient.getCuratorFramework().getConnectionStateListenable().addListener(onDisconnect); this.closeClient = true; this.securityNodeWatcher = null; collectionPropertiesZkStateReader = new CollectionPropertiesZkStateReader(this); assert ObjectReleaseTracker.track(this); } + private void onDisconnect(boolean sessionExpired) { + if (sessionExpired) { + zkSessionExpired.set(true); + } + } + + private void onReconnect() { + if (!zkSessionExpired.compareAndSet(true, false)) { + return; + } + try { + createClusterStateWatchersAndUpdate(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.warn("Interrupted", e); + } catch (Throwable e) { + log.error("An error has occurred while updating the cluster state", e); + } + } + /** * Forcibly refresh cluster state from ZK. Do this only to avoid race conditions because it's * expensive. diff --git a/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java b/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java index 9d6c83fa65af..4e3741f863a2 100644 --- a/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java +++ b/solr/solrj-zookeeper/src/test/org/apache/solr/common/cloud/TestOnReconnectSessionExpiry.java @@ -21,50 +21,40 @@ import org.apache.solr.SolrTestCase; import org.junit.Test; -/** SOLR-18298: Curator RECONNECTED is not the same as Solr 9 session-expiration reconnect. */ +/** Verifies the shared Curator listener adapters retain their general-purpose behavior. */ public class TestOnReconnectSessionExpiry extends SolrTestCase { @Test - public void testReconnectAfterSuspendDoesNotFire() { + public void testReconnectFiresForEveryReconnectedEvent() { AtomicInteger reconnects = new AtomicInteger(); OnReconnect listener = () -> reconnects.incrementAndGet(); listener.stateChanged(null, ConnectionState.SUSPENDED); listener.stateChanged(null, ConnectionState.RECONNECTED); - - assertEquals(0, reconnects.get()); - } - - @Test - public void testReconnectAfterLostDoesFire() { - AtomicInteger reconnects = new AtomicInteger(); - OnReconnect listener = () -> reconnects.incrementAndGet(); - listener.stateChanged(null, ConnectionState.LOST); listener.stateChanged(null, ConnectionState.RECONNECTED); - assertEquals(1, reconnects.get()); - } - - @Test - public void testReconnectWithoutPriorLostDoesNotFire() { - AtomicInteger reconnects = new AtomicInteger(); - OnReconnect listener = () -> reconnects.incrementAndGet(); - listener.stateChanged(null, ConnectionState.RECONNECTED); - - assertEquals(0, reconnects.get()); + assertEquals(3, reconnects.get()); } @Test - public void testDisconnectFiresOnlyOnLost() { - AtomicInteger disconnects = new AtomicInteger(); - OnDisconnect listener = sessionExpired -> disconnects.incrementAndGet(); + public void testDisconnectDistinguishesSuspensionFromSessionLoss() { + AtomicInteger suspensions = new AtomicInteger(); + AtomicInteger expirations = new AtomicInteger(); + OnDisconnect listener = + sessionExpired -> { + if (sessionExpired) { + expirations.incrementAndGet(); + } else { + suspensions.incrementAndGet(); + } + }; listener.stateChanged(null, ConnectionState.SUSPENDED); - assertEquals(0, disconnects.get()); - listener.stateChanged(null, ConnectionState.LOST); - assertEquals(1, disconnects.get()); + + assertEquals(1, suspensions.get()); + assertEquals(1, expirations.get()); } }