Skip to content

Commit eeddf25

Browse files
committed
fix: isolate suspect server state per client
Synchronize suspect cleanup with roster refresh and make stale cleanup idempotent.
1 parent a599227 commit eeddf25

7 files changed

Lines changed: 331 additions & 52 deletions

File tree

src/main/java/com/alipay/oceanbase/rpc/ObTableClient.java

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -944,7 +944,8 @@ public ObTable addTable(ObServerAddr addr){
944944
logger.info("server from response not exist in route cache, server ip {}, port {} , execute add Table.", addr.getIp(), addr.getSvrPort());
945945
ObTable obTable = new ObTable.Builder(addr.getIp(), addr.getSvrPort()) //
946946
.setLoginInfo(tenantName, userName, password, database, getClientType(runningMode)) //
947-
.setProperties(getProperties()).setObServerAddr(addr).build();
947+
.setProperties(getProperties()).setObServerAddr(addr)
948+
.setFailureHandler(tableRoute::reportObServerFailure).build();
948949
tableRoster.put(addr, obTable);
949950
return obTable;
950951
} catch (Exception e) {
@@ -983,14 +984,12 @@ public Row transformToRow(String tableName, Object[] rowkey) throws Exception {
983984
}
984985

985986
public void dealWithRpcTimeoutForSingleTablet(ObServerAddr addr, String tableName, long tabletId) throws Exception {
986-
RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr);
987-
RouteTableRefresher.addIntoSuspectIPs(suspectAddr);
987+
tableRoute.reportObServerFailure(addr);
988988
tableRoute.refreshPartitionLocation(tableName, tabletId, null);
989989
}
990990

991991
public void dealWithRpcTimeoutForBatchTablet(ObServerAddr addr, String tableName) throws Exception {
992-
RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr);
993-
RouteTableRefresher.addIntoSuspectIPs(suspectAddr);
992+
tableRoute.reportObServerFailure(addr);
994993
tableRoute.refreshTabletLocationBatch(tableName);
995994
}
996995

src/main/java/com/alipay/oceanbase/rpc/bolt/transport/ObTableConnection.java

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
import com.alipay.oceanbase.rpc.ObGlobal;
2121
import com.alipay.oceanbase.rpc.exception.*;
2222
import com.alipay.oceanbase.rpc.location.LocationUtil;
23-
import com.alipay.oceanbase.rpc.location.model.RouteTableRefresher;
2423
import com.alipay.oceanbase.rpc.protocol.payload.impl.login.ObTableLoginRequest;
2524
import com.alipay.oceanbase.rpc.protocol.payload.impl.login.ObTableLoginResult;
2625
import com.alipay.oceanbase.rpc.table.ObTable;
@@ -119,9 +118,7 @@ private boolean connect() throws Exception {
119118

120119
if (tries >= maxTryTimes) {
121120
if (!obTable.isOdpMode() && obTable.getObServerAddr() != null) {
122-
RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(
123-
obTable.getObServerAddr());
124-
RouteTableRefresher.addIntoSuspectIPs(suspectAddr);
121+
obTable.reportConnectionFailure();
125122
}
126123
LOGGER.warn("connect failed after max " + maxTryTimes + " tries "
127124
+ TraceUtil.formatIpPort(obTable));

src/main/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresher.java

Lines changed: 106 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -22,13 +22,13 @@
2222
import java.sql.Statement;
2323
import java.util.*;
2424
import java.util.concurrent.*;
25+
import java.util.concurrent.atomic.AtomicBoolean;
2526
import java.util.concurrent.locks.Lock;
2627
import java.util.concurrent.locks.ReentrantLock;
2728

2829
import com.alipay.oceanbase.rpc.ObTableClient;
2930
import com.alipay.oceanbase.rpc.exception.ObTableEntryRefreshException;
3031
import com.alipay.oceanbase.rpc.exception.ObTableTryLockTimeoutException;
31-
import com.alipay.oceanbase.rpc.exception.ObTableUnexpectedException;
3232
import com.alipay.oceanbase.rpc.location.LocationUtil;
3333
import com.alipay.oceanbase.rpc.table.ObTable;
3434
import org.slf4j.Logger;
@@ -49,11 +49,15 @@ public class RouteTableRefresher {
4949

5050
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2);
5151

52-
private final static ConcurrentHashMap<ObServerAddr, Lock> suspectLocks = new ConcurrentHashMap<>(); // ObServer -> access lock
52+
private final ConcurrentHashMap<ObServerAddr, Lock> suspectLocks = new ConcurrentHashMap<>(); // ObServer -> access lock
5353

54-
private final static ConcurrentHashMap<ObServerAddr, SuspectObServer> suspectServers = new ConcurrentHashMap<>(); // ObServer -> information structure
54+
private final ConcurrentHashMap<ObServerAddr, SuspectObServer> suspectServers = new ConcurrentHashMap<>(); // ObServer -> information structure
5555

56-
private final static HashMap<ObServerAddr, Long> serverLastAccessTimestamps = new HashMap<>(); // ObServer -> last access timestamp
56+
private final ConcurrentHashMap<ObServerAddr, Long> serverLastAccessTimestamps = new ConcurrentHashMap<>(); // ObServer -> last access timestamp
57+
58+
private final Set<ObServerAddr> activeServers = ConcurrentHashMap.newKeySet();
59+
60+
private final AtomicBoolean closed = new AtomicBoolean(false);
5761

5862
public RouteTableRefresher(ObTableClient tableClient, ObUserAuth sysUA) {
5963
this.tableClient = tableClient;
@@ -70,6 +74,9 @@ public void start() {
7074
}
7175

7276
public void close() {
77+
if (!closed.compareAndSet(false, true)) {
78+
return;
79+
}
7380
try {
7481
scheduler.shutdown();
7582
// wait at most 1 seconds to close the scheduler
@@ -79,6 +86,39 @@ public void close() {
7986
} catch (InterruptedException e) {
8087
logger.warn("scheduler await for terminate interrupted: {}.", e.getMessage());
8188
scheduler.shutdownNow();
89+
Thread.currentThread().interrupt();
90+
} finally {
91+
suspectServers.clear();
92+
suspectLocks.clear();
93+
serverLastAccessTimestamps.clear();
94+
activeServers.clear();
95+
}
96+
}
97+
98+
/**
99+
* Reconcile the keep-alive state with the latest authoritative tenant roster.
100+
* Servers removed from the roster must not remain in, or be re-added to, the suspect set.
101+
*/
102+
public void refreshActiveServers(Collection<ObServerAddr> servers) {
103+
if (closed.get()) {
104+
return;
105+
}
106+
Set<ObServerAddr> newServers = new HashSet<>();
107+
if (servers != null) {
108+
newServers.addAll(servers);
109+
}
110+
activeServers.retainAll(newServers);
111+
activeServers.addAll(newServers);
112+
113+
for (ObServerAddr addr : suspectServers.keySet()) {
114+
if (!newServers.contains(addr)) {
115+
removeFromSuspectIPs(addr);
116+
}
117+
}
118+
for (ObServerAddr addr : serverLastAccessTimestamps.keySet()) {
119+
if (!newServers.contains(addr)) {
120+
serverLastAccessTimestamps.remove(addr);
121+
}
82122
}
83123
}
84124

@@ -127,6 +167,10 @@ private void doRsListCheck() {
127167
private void doCheckAliveTask() {
128168
for (Map.Entry<ObServerAddr, SuspectObServer> entry : suspectServers.entrySet()) {
129169
try {
170+
if (!activeServers.contains(entry.getKey())) {
171+
removeFromSuspectIPs(entry.getKey());
172+
continue;
173+
}
130174
checkAlive(entry.getKey());
131175
} catch (Exception e) {
132176
// silence resolving
@@ -161,7 +205,7 @@ private void checkAlive(ObServerAddr addr) {
161205
if (t instanceof SQLException) {
162206
// occurred during query
163207
calcFailureOrClearCache(addr);
164-
} if (t instanceof ObTableEntryRefreshException) {
208+
} else if (t instanceof ObTableEntryRefreshException) {
165209
// occurred during connection construction
166210
ObTableEntryRefreshException e = (ObTableEntryRefreshException) t;
167211
if (e.isConnectInactive()) {
@@ -192,11 +236,22 @@ private void checkAlive(ObServerAddr addr) {
192236
}
193237
}
194238

195-
public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedException {
239+
public void addIntoSuspectIPs(ObServerAddr addr) {
240+
if (addr == null) {
241+
return;
242+
}
243+
addIntoSuspectIPs(new SuspectObServer(addr));
244+
}
245+
246+
private void addIntoSuspectIPs(SuspectObServer server) {
196247
if (server == null || server.getAddr() == null) {
197248
return;
198249
}
199250
ObServerAddr addr = server.getAddr();
251+
if (closed.get() || !activeServers.contains(addr)) {
252+
logger.debug("ignore suspect report for inactive server: {}", addr);
253+
return;
254+
}
200255
if (suspectServers.get(addr) != null) {
201256
// already in the list, directly return
202257
return;
@@ -218,6 +273,9 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE
218273
// already in the list, directly break
219274
break;
220275
}
276+
if (closed.get() || !activeServers.contains(addr)) {
277+
break;
278+
}
221279
Long lastServerAccessTs = serverLastAccessTimestamps.get(addr);
222280
if (lastServerAccessTs != null) {
223281
long interval = System.currentTimeMillis() - lastServerAccessTs;
@@ -235,6 +293,11 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE
235293
++retryTimes;
236294
logger.warn("wait to try lock to timeout 1s when add observer into suspect ips, server: {}, tryTimes: {}",
237295
addr.toString(), retryTimes, e);
296+
} catch (InterruptedException e) {
297+
Thread.currentThread().interrupt();
298+
logger.debug("interrupted while adding observer into suspect ips, server: {}",
299+
addr);
300+
break;
238301
}
239302
} // end while
240303
} finally {
@@ -247,43 +310,33 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE
247310
private void removeFromSuspectIPs(ObServerAddr addr) {
248311
Lock lock = suspectLocks.get(addr);
249312
if (lock == null) {
250-
// lock must have been added before remove
251-
throw new ObTableUnexpectedException(String.format("ObServer [%s:%d] need to be add into suspect ips before remove",
252-
addr.getIp(), addr.getSvrPort()));
313+
suspectServers.remove(addr);
314+
logger.debug("suspect server has already been removed: {}", addr);
315+
return;
253316
}
254317
boolean acquired = false;
255318
try {
256-
int retryTimes = 0;
257-
while (true) {
258-
try {
259-
acquired = lock.tryLock(1, TimeUnit.SECONDS);
260-
if (!acquired) {
261-
throw new ObTableTryLockTimeoutException("try to get suspect server lock timeout, timeout: 1s");
262-
}
263-
// no need to remove lock
264-
SuspectObServer server = suspectServers.remove(addr);
265-
if (server != null) {
266-
int failure = server.getFailure();
267-
if (failure < failureLimit) {
268-
ObTable obTable = tableClient.getTableRoute().getTableRoster().getTable(addr);
269-
if (obTable != null && !obTable.isValid()) {
270-
obTable.setValid();
271-
}
272-
}
319+
acquired = lock.tryLock(1, TimeUnit.SECONDS);
320+
if (!acquired) {
321+
logger.debug("defer suspect removal because lock is busy, server: {}", addr);
322+
return;
323+
}
324+
// Keep the lock and cooldown entry until this refresher closes. This prevents a
325+
// concurrent add from using a different lock and preserves the existing cooldown.
326+
SuspectObServer server = suspectServers.remove(addr);
327+
if (server != null) {
328+
int failure = server.getFailure();
329+
if (failure < failureLimit && activeServers.contains(addr)) {
330+
ObTable obTable = tableClient.getTableRoute().getTableRoster().getTable(addr);
331+
if (obTable != null && !obTable.isValid()) {
332+
obTable.setValid();
273333
}
274-
logger.debug("removed server from suspect list: {}", addr);
275-
break;
276-
} catch (ObTableTryLockTimeoutException e) {
277-
// if try lock timeout, need to retry
278-
++retryTimes;
279-
logger.warn("wait to try lock to timeout when add observer into suspect ips, server: {}, tryTimes: {}",
280-
addr.toString(), retryTimes, e);
281-
} catch (InterruptedException e) {
282-
// do not throw exception to user layer
283-
// next background task will continue to remove it
284-
logger.warn("waiting to get lock while interrupted by other threads", e);
285334
}
286335
}
336+
logger.debug("removed server from suspect list: {}", addr);
337+
} catch (InterruptedException e) {
338+
Thread.currentThread().interrupt();
339+
logger.debug("interrupted while removing observer from suspect ips, server: {}", addr);
287340
} finally {
288341
if (acquired) {
289342
lock.unlock();
@@ -294,6 +347,10 @@ private void removeFromSuspectIPs(ObServerAddr addr) {
294347
private void calcFailureOrClearCache(ObServerAddr addr) {
295348
TableRoute tableRoute = tableClient.getTableRoute();
296349
SuspectObServer server = suspectServers.get(addr);
350+
if (server == null) {
351+
logger.debug("skip failure calculation for removed suspect server: {}", addr);
352+
return;
353+
}
297354
server.incrementFailure();
298355
int failure = server.getFailure();
299356
if (failure >= failureLimit) {
@@ -304,6 +361,18 @@ private void calcFailureOrClearCache(ObServerAddr addr) {
304361
addr, failure);
305362
}
306363

364+
int suspectServerCount() {
365+
return suspectServers.size();
366+
}
367+
368+
boolean containsSuspectServer(ObServerAddr addr) {
369+
return suspectServers.containsKey(addr);
370+
}
371+
372+
int lastAccessTimestampCount() {
373+
return serverLastAccessTimestamps.size();
374+
}
375+
307376
public static class SuspectObServer {
308377
private final ObServerAddr addr;
309378
private final long accessTimestamp;

src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoster.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import java.util.*;
2121
import java.util.concurrent.ConcurrentHashMap;
22+
import java.util.function.Consumer;
2223

2324
import com.alipay.oceanbase.rpc.ObTableClient;
2425
import com.alipay.oceanbase.rpc.exception.ObTableCloseException;
@@ -36,6 +37,7 @@ public class TableRoster {
3637
private Properties properties = new Properties();
3738
private Map<String, Object> tableConfigs = new HashMap<>();
3839
private ObTableClientType clientType;
40+
private Consumer<ObServerAddr> failureHandler;
3941
/*
4042
* ServerAddr(all) -> ObTableConnection
4143
*/
@@ -65,6 +67,9 @@ public void setProperties(Properties properties) {
6567
public void setTableConfigs(Map<String, Object> tableConfigs) {
6668
this.tableConfigs = tableConfigs;
6769
}
70+
public void setFailureHandler(Consumer<ObServerAddr> failureHandler) {
71+
this.failureHandler = failureHandler;
72+
}
6873
public ObTable getTable(ObServerAddr addr) {
6974
return tables.get(addr);
7075
}
@@ -101,7 +106,8 @@ public List<ObServerAddr> refreshTablesAndGetNewServers(List<ReplicaLocation> ne
101106

102107
ObTable obTable = new ObTable.Builder(addr.getIp(), addr.getSvrPort()) //
103108
.setLoginInfo(tenantName, userName, password, database, clientType) //
104-
.setProperties(properties).setConfigs(tableConfigs).setObServerAddr(addr).build();
109+
.setProperties(properties).setConfigs(tableConfigs).setObServerAddr(addr)
110+
.setFailureHandler(failureHandler).build();
105111
ObTable oldObTable = tables.putIfAbsent(addr, obTable);
106112
logger.warn("add new table addr, {}", addr.toString());
107113
if (oldObTable != null) { // maybe create two ob table concurrently, close current ob table

src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoute.java

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -64,7 +64,7 @@ public class TableRoute {
6464
private IndexLocations indexLocations = null; // global index location
6565
private TableGroupCache tableGroupCache = null;
6666
private OdpInfo odpInfo = null;
67-
private RouteTableRefresher routeRefresher = null;
67+
private volatile RouteTableRefresher routeRefresher = null;
6868
private long lastRefreshMetadataTimestamp = -1;
6969
public final Lock refreshTableRosterLock = new ReentrantLock();
7070

@@ -348,7 +348,8 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized,
348348
tableClient.getPassword(), tableClient.getDatabase(),
349349
tableClient.getClientType(runningMode))
350350
.setProperties(tableClient.getProperties())
351-
.setConfigs(tableClient.getTableConfigs()).setObServerAddr(addr).build();
351+
.setConfigs(tableClient.getTableConfigs()).setObServerAddr(addr)
352+
.setFailureHandler(this::reportObServerFailure).build();
352353
addr2Table.put(addr, obTable);
353354
servers.add(addr);
354355
} catch (Exception e) {
@@ -363,6 +364,7 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized,
363364
tableClient.getUserName(), tableClient.getPassword(), tableClient.getDatabase(),
364365
tableClient.getClientType(runningMode), tableClient.getProperties(),
365366
tableClient.getTableConfigs());
367+
this.tableRoster.setFailureHandler(this::reportObServerFailure);
366368
this.tableRoster.setTables(addr2Table);
367369
this.serverRoster.reset(servers);
368370

@@ -410,9 +412,17 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized,
410412

411413
public void launchRouteRefresher() {
412414
routeRefresher = new RouteTableRefresher(tableClient, sysUA);
415+
routeRefresher.refreshActiveServers(serverRoster.getMembers());
413416
routeRefresher.start();
414417
}
415418

419+
public void reportObServerFailure(ObServerAddr addr) {
420+
RouteTableRefresher refresher = routeRefresher;
421+
if (refresher != null) {
422+
refresher.addIntoSuspectIPs(addr);
423+
}
424+
}
425+
416426
public void removeObServer(ObServerAddr addr) {
417427
logger.debug("remove useless table addr, {}", addr.toString());
418428
ConcurrentHashMap<ObServerAddr, ObTable> tables = this.tableRoster.getTables();
@@ -480,6 +490,10 @@ public void refreshRosterByRsList(List<ObServerAddr> newRsList) throws Exception
480490
// update new ob table and get new server address
481491
List<ObServerAddr> servers = tableRoster.refreshTablesAndGetNewServers(replicaLocations);
482492
serverRoster.reset(servers);
493+
RouteTableRefresher refresher = routeRefresher;
494+
if (refresher != null) {
495+
refresher.refreshActiveServers(servers);
496+
}
483497

484498
// 2. Get Server LDC info for weak read consistency.
485499
success = false;

0 commit comments

Comments
 (0)