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());
+
+ }
}