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

kturner pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/accumulo.git


The following commit(s) were added to refs/heads/main by this push:
     new dff80ffaca handles same host+port in different resource groups in 
LiveTServerSet (#6302)
dff80ffaca is described below

commit dff80ffacaa4c87ecbb4acacc3a2f4d1d644cfa8
Author: Keith Turner <[email protected]>
AuthorDate: Wed Apr 8 07:49:23 2026 -0700

    handles same host+port in different resource groups in LiveTServerSet 
(#6302)
    
    The same host+port in different resource groups in ZK would cause
    LiveTServerSet to continually add and remove the host port. Changed an
    internal map that keyed on host+port to instead key on the ZK path to
    fix this.
    
    fixes #6298
---
 .../accumulo/server/manager/LiveTServerSet.java    |  57 ++++++----
 .../server/manager/LiveTServerSetTest.java         | 125 ++++++++++++++++++++-
 2 files changed, 152 insertions(+), 30 deletions(-)

diff --git 
a/server/base/src/main/java/org/apache/accumulo/server/manager/LiveTServerSet.java
 
b/server/base/src/main/java/org/apache/accumulo/server/manager/LiveTServerSet.java
index 8d7f5ee653..3e86c0441d 100644
--- 
a/server/base/src/main/java/org/apache/accumulo/server/manager/LiveTServerSet.java
+++ 
b/server/base/src/main/java/org/apache/accumulo/server/manager/LiveTServerSet.java
@@ -78,7 +78,7 @@ public class LiveTServerSet implements ZooCacheWatcher {
 
   private static final Logger log = 
LoggerFactory.getLogger(LiveTServerSet.class);
 
-  private final AtomicReference<Listener> cback;
+  protected final AtomicReference<Listener> cback;
   private final ServerContext context;
 
   public class TServerConnection {
@@ -204,7 +204,7 @@ public class LiveTServerSet implements ZooCacheWatcher {
 
   // The set of active tservers with locks, indexed by their name in 
zookeeper. When the contents of
   // this map are modified, tServersSnapshot should be set to null.
-  private final Map<String,TServerInfo> current = new HashMap<>();
+  private final Map<ServiceLockPath,TServerInfo> current = new HashMap<>();
 
   private LiveTServersSnapshot tServersSnapshot = null;
 
@@ -233,12 +233,17 @@ public class LiveTServerSet implements ZooCacheWatcher {
         .scheduleWithFixedDelay(this::scanServers, 5000, 5000, 
TimeUnit.MILLISECONDS));
   }
 
+  @VisibleForTesting
+  protected Set<ServiceLockPath> getTserverPaths() {
+    return context.getServerPaths().getTabletServer(ResourceGroupPredicate.ANY,
+        AddressSelector.all(), false);
+  }
+
   public synchronized void scanServers() {
     try {
       final Set<TServerInstance> updates = new HashSet<>();
       final Set<TServerInstance> doomed = new HashSet<>();
-      final Set<ServiceLockPath> tservers = context.getServerPaths()
-          .getTabletServer(ResourceGroupPredicate.ANY, AddressSelector.all(), 
false);
+      final Set<ServiceLockPath> tservers = getTserverPaths();
 
       locklessServers.keySet().retainAll(tservers);
 
@@ -262,6 +267,11 @@ public class LiveTServerSet implements ZooCacheWatcher {
     }
   }
 
+  @VisibleForTesting
+  protected Optional<ServiceLockData> getLockData(ServiceLockPath tserverPath, 
ZcStat stat) {
+    return ServiceLock.getLockData(context.getZooCache(), tserverPath, stat);
+  }
+
   private synchronized void checkServer(final Set<TServerInstance> updates,
       final Set<TServerInstance> doomed, final ServiceLockPath tserverPath)
       throws InterruptedException, KeeperException {
@@ -269,33 +279,32 @@ public class LiveTServerSet implements ZooCacheWatcher {
     // invalidate the snapshot forcing it to be recomputed the next time its 
requested
     tServersSnapshot = null;
 
-    final TServerInfo info = current.get(tserverPath.getServer());
+    final TServerInfo info = current.get(tserverPath);
 
     ZcStat stat = new ZcStat();
-    Optional<ServiceLockData> sld =
-        ServiceLock.getLockData(context.getZooCache(), tserverPath, stat);
+    Optional<ServiceLockData> sld = getLockData(tserverPath, stat);
 
     if (sld.isEmpty()) {
-      log.trace("lock does not exist for server: {}", tserverPath.getServer());
+      log.trace("lock does not exist for server: {}", tserverPath);
       if (info != null) {
         doomed.add(info.instance);
-        current.remove(tserverPath.getServer());
-        log.trace("removed {} from current set and adding to doomed list", 
tserverPath.getServer());
+        current.remove(tserverPath);
+        log.trace("removed {} from current set and adding to doomed list", 
tserverPath);
       }
 
       Long firstSeen = locklessServers.get(tserverPath);
       if (firstSeen == null) {
         locklessServers.put(tserverPath, System.currentTimeMillis());
-        log.trace("first seen, added {} to list of lockless servers", 
tserverPath.getServer());
+        log.trace("first seen, added {} to list of lockless servers", 
tserverPath);
       } else if (System.currentTimeMillis() - firstSeen > 
MINUTES.toMillis(10)) {
         deleteServerNode(tserverPath.toString());
         locklessServers.remove(tserverPath);
         log.trace(
             "deleted zookeeper node for server: {}, has been without lock for 
over 10 minutes",
-            tserverPath.getServer());
+            tserverPath);
       }
     } else {
-      log.trace("Lock exists for server: {}, adding to current set", 
tserverPath.getServer());
+      log.trace("Lock exists for server: {}, adding to current set", 
tserverPath);
       locklessServers.remove(tserverPath);
       HostAndPort address = 
sld.orElseThrow().getAddress(ServiceLockData.ThriftService.TSERV);
       ResourceGroupId resourceGroup =
@@ -306,13 +315,13 @@ public class LiveTServerSet implements ZooCacheWatcher {
         updates.add(instance);
         TServerInfo tServerInfo =
             new TServerInfo(instance, new TServerConnection(address), 
resourceGroup);
-        current.put(tserverPath.getServer(), tServerInfo);
+        current.put(tserverPath, tServerInfo);
       } else if (!info.instance.equals(instance)) {
         doomed.add(info.instance);
         updates.add(instance);
         TServerInfo tServerInfo =
             new TServerInfo(instance, new TServerConnection(address), 
resourceGroup);
-        current.put(tserverPath.getServer(), tServerInfo);
+        current.put(tserverPath, tServerInfo);
       }
     }
   }
@@ -459,7 +468,7 @@ public class LiveTServerSet implements ZooCacheWatcher {
     return find(current, tabletServer);
   }
 
-  static TServerInstance find(Map<String,TServerInfo> servers, String 
tabletServer) {
+  static TServerInstance find(Map<ServiceLockPath,TServerInfo> servers, String 
tabletServer) {
     HostAndPort addr;
     String sessionId = null;
     if (tabletServer.charAt(tabletServer.length() - 1) == ']') {
@@ -473,11 +482,11 @@ public class LiveTServerSet implements ZooCacheWatcher {
     } else {
       addr = AddressUtil.parseAddress(tabletServer);
     }
-    for (Entry<String,TServerInfo> entry : servers.entrySet()) {
-      if (entry.getValue().instance.getHostAndPort().equals(addr)) {
+    for (TServerInfo tServerInfo : servers.values()) {
+      if (tServerInfo.instance.getHostAndPort().equals(addr)) {
         // Return the instance if we have no desired session ID, or we match 
the desired session ID
-        if (sessionId == null || 
sessionId.equals(entry.getValue().instance.getSession())) {
-          return entry.getValue().instance;
+        if (sessionId == null || 
sessionId.equals(tServerInfo.instance.getSession())) {
+          return tServerInfo.instance;
         }
       }
     }
@@ -491,17 +500,17 @@ public class LiveTServerSet implements ZooCacheWatcher {
 
     Optional<ResourceGroupId> resourceGroup = Optional.empty();
     Optional<HostAndPort> address = Optional.empty();
-    for (Entry<String,TServerInfo> entry : current.entrySet()) {
+    for (Entry<ServiceLockPath,TServerInfo> entry : current.entrySet()) {
       if (entry.getValue().instance.equals(server)) {
-        address = Optional.of(HostAndPort.fromString(entry.getKey()));
+        address = 
Optional.of(HostAndPort.fromString(entry.getKey().getServer()));
         resourceGroup = Optional.of(entry.getValue().resourceGroup);
+        current.remove(entry.getKey());
         break;
       }
     }
     if (resourceGroup.isEmpty() || address.isEmpty()) {
       return;
     }
-    current.remove(address.orElseThrow().toString());
 
     ResourceGroupPredicate rgPredicate = resourceGroup.map(rg -> {
       ResourceGroupPredicate rgp = rg2 -> rg.equals(rg2);
@@ -511,7 +520,7 @@ public class LiveTServerSet implements ZooCacheWatcher {
         address.map(AddressSelector::exact).orElse(AddressSelector.all());
     Set<ServiceLockPath> paths =
         context.getServerPaths().getTabletServer(rgPredicate, addrPredicate, 
false);
-    if (paths.isEmpty() || paths.size() > 1) {
+    if (paths.size() != 1) {
       log.error("Zero or many zookeeper entries match input arguments.");
     } else {
       ServiceLockPath slp = paths.iterator().next();
diff --git 
a/server/base/src/test/java/org/apache/accumulo/server/manager/LiveTServerSetTest.java
 
b/server/base/src/test/java/org/apache/accumulo/server/manager/LiveTServerSetTest.java
index 15e9d5e3ac..837b5f7fac 100644
--- 
a/server/base/src/test/java/org/apache/accumulo/server/manager/LiveTServerSetTest.java
+++ 
b/server/base/src/test/java/org/apache/accumulo/server/manager/LiveTServerSetTest.java
@@ -18,32 +18,49 @@
  */
 package org.apache.accumulo.server.manager;
 
+import static org.easymock.EasyMock.createMock;
+import static org.easymock.EasyMock.createStrictMock;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
 
 import java.util.HashMap;
+import java.util.HashSet;
 import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.UUID;
 
 import org.apache.accumulo.core.data.ResourceGroupId;
+import org.apache.accumulo.core.lock.ServiceLockData;
+import org.apache.accumulo.core.lock.ServiceLockData.ThriftService;
+import org.apache.accumulo.core.lock.ServiceLockPaths;
+import org.apache.accumulo.core.lock.ServiceLockPaths.ServiceLockPath;
 import org.apache.accumulo.core.metadata.TServerInstance;
+import org.apache.accumulo.core.zookeeper.ZcStat;
+import org.apache.accumulo.core.zookeeper.ZooCache;
+import org.apache.accumulo.server.ServerContext;
 import org.apache.accumulo.server.manager.LiveTServerSet.TServerConnection;
 import org.apache.accumulo.server.manager.LiveTServerSet.TServerInfo;
-import org.easymock.EasyMock;
 import org.junit.jupiter.api.Test;
 
+import com.google.common.annotations.VisibleForTesting;
 import com.google.common.net.HostAndPort;
 
 public class LiveTServerSetTest {
 
   @Test
   public void testSessionIds() {
-    Map<String,TServerInfo> servers = new HashMap<>();
-    TServerConnection mockConn = EasyMock.createMock(TServerConnection.class);
+    Map<ServiceLockPath,TServerInfo> servers = new HashMap<>();
+    TServerConnection mockConn = createMock(TServerConnection.class);
+
+    var hostPort = HostAndPort.fromParts("localhost", 1234);
 
     TServerInfo server1 =
-        new TServerInfo(new TServerInstance(HostAndPort.fromParts("localhost", 
1234), "5555"),
-            mockConn, ResourceGroupId.DEFAULT);
-    servers.put("server1", server1);
+        new TServerInfo(new TServerInstance(hostPort, "5555"), mockConn, 
ResourceGroupId.DEFAULT);
+
+    ServiceLockPaths lockPaths = new 
ServiceLockPaths(createMock(ZooCache.class));
+
+    servers.put(lockPaths.createTabletServerPath(ResourceGroupId.DEFAULT, 
hostPort), server1);
 
     assertEquals(server1.instance, LiveTServerSet.find(servers, 
"localhost:1234"));
     assertNull(LiveTServerSet.find(servers, "localhost:4321"));
@@ -51,4 +68,100 @@ public class LiveTServerSetTest {
     assertNull(LiveTServerSet.find(servers, "localhost:1234[55755]"));
   }
 
+  record LockInfo(ServiceLockData sld, long sessionId) {
+  }
+
+  static class TestLiveTserverSet extends LiveTServerSet {
+    private final Set<ServiceLockPath> paths;
+    private final Map<ServiceLockPath,LockInfo> locks;
+
+    public TestLiveTserverSet(Set<ServiceLockPath> paths, 
Map<ServiceLockPath,LockInfo> locks,
+        Listener listener) {
+      super(createStrictMock(ServerContext.class));
+      this.paths = paths;
+      this.locks = locks;
+      this.cback.set(listener);
+    }
+
+    @VisibleForTesting
+    protected Set<ServiceLockPath> getTserverPaths() {
+      return paths;
+    }
+
+    @VisibleForTesting
+    protected Optional<ServiceLockData> getLockData(ServiceLockPath 
tserverPath, ZcStat stat) {
+      var lockInfo = locks.get(tserverPath);
+      if (lockInfo == null) {
+        return Optional.empty();
+      }
+
+      stat.setEphemeralOwner(lockInfo.sessionId);
+
+      return Optional.of(lockInfo.sld);
+    }
+  }
+
+  @Test
+  public void testSameHostPortInDiffResourceGroups() {
+    // This tests two tservers with the same host port in different resource 
groups where only one
+    // has a lock. There was a bug where LiveTserverSet would add and remove 
that host port from its
+    // set and pass both add/remove to the callback listener.
+
+    ServiceLockPaths lockPaths = new 
ServiceLockPaths(createMock(ZooCache.class));
+
+    var hp1 = HostAndPort.fromParts("host1", 9800);
+    var g1 = ResourceGroupId.of("tg1");
+    var path1 = lockPaths.createTabletServerPath(ResourceGroupId.DEFAULT, hp1);
+
+    var path2 = lockPaths.createTabletServerPath(g1, hp1);
+
+    var paths = new HashSet<ServiceLockPath>();
+    paths.add(path1);
+    paths.add(path2);
+
+    var locks = new HashMap<ServiceLockPath,LockInfo>();
+    locks.put(path2, new LockInfo(
+        new ServiceLockData(UUID.randomUUID(), hp1.toString(), 
ThriftService.TSERV, g1), 123456));
+
+    var deletedSeen = new HashSet<TServerInstance>();
+    var addedSeen = new HashSet<TServerInstance>();
+    LiveTServerSet tservers = new TestLiveTserverSet(paths, locks, ((current, 
deleted, added) -> {
+      deletedSeen.addAll(deleted);
+      addedSeen.addAll(added);
+    }));
+    tservers.scanServers();
+
+    var expected1 = Set.of(new TServerInstance(hp1, 123456));
+    assertEquals(expected1, addedSeen);
+    assertEquals(Set.of(), deletedSeen);
+    assertEquals(expected1, tservers.getSnapshot().getTservers());
+    assertEquals(Map.of(g1, expected1), 
tservers.getSnapshot().getTserverGroups());
+
+    // change which tserver has the lock
+    locks.clear();
+    locks.put(path1, new LockInfo(new ServiceLockData(UUID.randomUUID(), 
hp1.toString(),
+        ThriftService.TSERV, ResourceGroupId.DEFAULT), 654321));
+
+    addedSeen.clear();
+    tservers.scanServers();
+
+    var expected2 = Set.of(new TServerInstance(hp1, 654321));
+    assertEquals(expected2, addedSeen);
+    assertEquals(expected1, deletedSeen);
+    assertEquals(expected2, tservers.getSnapshot().getTservers());
+    assertEquals(Map.of(ResourceGroupId.DEFAULT, expected2),
+        tservers.getSnapshot().getTserverGroups());
+
+    // test when nothing has changed since last scan
+    addedSeen.clear();
+    deletedSeen.clear();
+    tservers.scanServers();
+    // Should not see any add/removes
+    assertEquals(Set.of(), addedSeen);
+    assertEquals(Set.of(), deletedSeen);
+    assertEquals(expected2, tservers.getSnapshot().getTservers());
+    assertEquals(Map.of(ResourceGroupId.DEFAULT, expected2),
+        tservers.getSnapshot().getTserverGroups());
+
+  }
 }

Reply via email to