Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 4 additions & 5 deletions src/main/java/com/alipay/oceanbase/rpc/ObTableClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -944,7 +944,8 @@ public ObTable addTable(ObServerAddr addr){
logger.info("server from response not exist in route cache, server ip {}, port {} , execute add Table.", addr.getIp(), addr.getSvrPort());
ObTable obTable = new ObTable.Builder(addr.getIp(), addr.getSvrPort()) //
.setLoginInfo(tenantName, userName, password, database, getClientType(runningMode)) //
.setProperties(getProperties()).setObServerAddr(addr).build();
.setProperties(getProperties()).setObServerAddr(addr)
.setFailureHandler(tableRoute::reportObServerFailure).build();
tableRoster.put(addr, obTable);
return obTable;
} catch (Exception e) {
Expand Down Expand Up @@ -983,14 +984,12 @@ public Row transformToRow(String tableName, Object[] rowkey) throws Exception {
}

public void dealWithRpcTimeoutForSingleTablet(ObServerAddr addr, String tableName, long tabletId) throws Exception {
RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr);
RouteTableRefresher.addIntoSuspectIPs(suspectAddr);
tableRoute.reportObServerFailure(addr);
tableRoute.refreshPartitionLocation(tableName, tabletId, null);
}

public void dealWithRpcTimeoutForBatchTablet(ObServerAddr addr, String tableName) throws Exception {
RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr);
RouteTableRefresher.addIntoSuspectIPs(suspectAddr);
tableRoute.reportObServerFailure(addr);
tableRoute.refreshTabletLocationBatch(tableName);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@
import com.alipay.oceanbase.rpc.ObGlobal;
import com.alipay.oceanbase.rpc.exception.*;
import com.alipay.oceanbase.rpc.location.LocationUtil;
import com.alipay.oceanbase.rpc.location.model.RouteTableRefresher;
import com.alipay.oceanbase.rpc.protocol.payload.impl.login.ObTableLoginRequest;
import com.alipay.oceanbase.rpc.protocol.payload.impl.login.ObTableLoginResult;
import com.alipay.oceanbase.rpc.table.ObTable;
Expand Down Expand Up @@ -119,9 +118,7 @@ private boolean connect() throws Exception {

if (tries >= maxTryTimes) {
if (!obTable.isOdpMode() && obTable.getObServerAddr() != null) {
RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(
obTable.getObServerAddr());
RouteTableRefresher.addIntoSuspectIPs(suspectAddr);
obTable.reportConnectionFailure();
}
LOGGER.warn("connect failed after max " + maxTryTimes + " tries "
+ TraceUtil.formatIpPort(obTable));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,13 +22,13 @@
import java.sql.Statement;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;

import com.alipay.oceanbase.rpc.ObTableClient;
import com.alipay.oceanbase.rpc.exception.ObTableEntryRefreshException;
import com.alipay.oceanbase.rpc.exception.ObTableTryLockTimeoutException;
import com.alipay.oceanbase.rpc.exception.ObTableUnexpectedException;
import com.alipay.oceanbase.rpc.location.LocationUtil;
import com.alipay.oceanbase.rpc.table.ObTable;
import org.slf4j.Logger;
Expand All @@ -49,11 +49,15 @@ public class RouteTableRefresher {

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

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

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

private final static HashMap<ObServerAddr, Long> serverLastAccessTimestamps = new HashMap<>(); // ObServer -> last access timestamp
private final ConcurrentHashMap<ObServerAddr, Long> serverLastAccessTimestamps = new ConcurrentHashMap<>(); // ObServer -> last access timestamp

private final Set<ObServerAddr> activeServers = ConcurrentHashMap.newKeySet();

private final AtomicBoolean closed = new AtomicBoolean(false);

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

public void close() {
if (!closed.compareAndSet(false, true)) {
return;
}
try {
scheduler.shutdown();
// wait at most 1 seconds to close the scheduler
Expand All @@ -79,6 +86,39 @@ public void close() {
} catch (InterruptedException e) {
logger.warn("scheduler await for terminate interrupted: {}.", e.getMessage());
scheduler.shutdownNow();
Thread.currentThread().interrupt();
} finally {
suspectServers.clear();
suspectLocks.clear();
serverLastAccessTimestamps.clear();
activeServers.clear();
}
}

/**
* Reconcile the keep-alive state with the latest authoritative tenant roster.
* Servers removed from the roster must not remain in, or be re-added to, the suspect set.
*/
public void refreshActiveServers(Collection<ObServerAddr> servers) {
if (closed.get()) {
return;
}
Set<ObServerAddr> newServers = new HashSet<>();
if (servers != null) {
newServers.addAll(servers);
}
activeServers.retainAll(newServers);
activeServers.addAll(newServers);

for (ObServerAddr addr : suspectServers.keySet()) {
if (!newServers.contains(addr)) {
removeFromSuspectIPs(addr);
}
}
for (ObServerAddr addr : serverLastAccessTimestamps.keySet()) {
if (!newServers.contains(addr)) {
serverLastAccessTimestamps.remove(addr);
}
}
}

Expand Down Expand Up @@ -127,6 +167,10 @@ private void doRsListCheck() {
private void doCheckAliveTask() {
for (Map.Entry<ObServerAddr, SuspectObServer> entry : suspectServers.entrySet()) {
try {
if (!activeServers.contains(entry.getKey())) {
removeFromSuspectIPs(entry.getKey());
continue;
}
checkAlive(entry.getKey());
} catch (Exception e) {
// silence resolving
Expand Down Expand Up @@ -161,7 +205,7 @@ private void checkAlive(ObServerAddr addr) {
if (t instanceof SQLException) {
// occurred during query
calcFailureOrClearCache(addr);
} if (t instanceof ObTableEntryRefreshException) {
} else if (t instanceof ObTableEntryRefreshException) {
// occurred during connection construction
ObTableEntryRefreshException e = (ObTableEntryRefreshException) t;
if (e.isConnectInactive()) {
Expand Down Expand Up @@ -192,11 +236,22 @@ private void checkAlive(ObServerAddr addr) {
}
}

public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedException {
public void addIntoSuspectIPs(ObServerAddr addr) {
if (addr == null) {
return;
}
addIntoSuspectIPs(new SuspectObServer(addr));
}

private void addIntoSuspectIPs(SuspectObServer server) {
if (server == null || server.getAddr() == null) {
return;
}
ObServerAddr addr = server.getAddr();
if (closed.get() || !activeServers.contains(addr)) {
logger.debug("ignore suspect report for inactive server: {}", addr);
return;
}
if (suspectServers.get(addr) != null) {
// already in the list, directly return
return;
Expand All @@ -218,6 +273,9 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE
// already in the list, directly break
break;
}
if (closed.get() || !activeServers.contains(addr)) {
break;
}
Long lastServerAccessTs = serverLastAccessTimestamps.get(addr);
if (lastServerAccessTs != null) {
long interval = System.currentTimeMillis() - lastServerAccessTs;
Expand All @@ -235,6 +293,11 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE
++retryTimes;
logger.warn("wait to try lock to timeout 1s when add observer into suspect ips, server: {}, tryTimes: {}",
addr.toString(), retryTimes, e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
logger.debug("interrupted while adding observer into suspect ips, server: {}",
addr);
break;
}
} // end while
} finally {
Expand All @@ -247,43 +310,33 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE
private void removeFromSuspectIPs(ObServerAddr addr) {
Lock lock = suspectLocks.get(addr);
if (lock == null) {
// lock must have been added before remove
throw new ObTableUnexpectedException(String.format("ObServer [%s:%d] need to be add into suspect ips before remove",
addr.getIp(), addr.getSvrPort()));
suspectServers.remove(addr);
logger.debug("suspect server has already been removed: {}", addr);
return;
}
boolean acquired = false;
try {
int retryTimes = 0;
while (true) {
try {
acquired = lock.tryLock(1, TimeUnit.SECONDS);
if (!acquired) {
throw new ObTableTryLockTimeoutException("try to get suspect server lock timeout, timeout: 1s");
}
// no need to remove lock
SuspectObServer server = suspectServers.remove(addr);
if (server != null) {
int failure = server.getFailure();
if (failure < failureLimit) {
ObTable obTable = tableClient.getTableRoute().getTableRoster().getTable(addr);
if (obTable != null && !obTable.isValid()) {
obTable.setValid();
}
}
acquired = lock.tryLock(1, TimeUnit.SECONDS);
if (!acquired) {
logger.debug("defer suspect removal because lock is busy, server: {}", addr);
return;
}
// Keep the lock and cooldown entry until this refresher closes. This prevents a
// concurrent add from using a different lock and preserves the existing cooldown.
SuspectObServer server = suspectServers.remove(addr);
if (server != null) {
int failure = server.getFailure();
if (failure < failureLimit && activeServers.contains(addr)) {
ObTable obTable = tableClient.getTableRoute().getTableRoster().getTable(addr);
if (obTable != null && !obTable.isValid()) {
obTable.setValid();
}
logger.debug("removed server from suspect list: {}", addr);
break;
} catch (ObTableTryLockTimeoutException e) {
// if try lock timeout, need to retry
++retryTimes;
logger.warn("wait to try lock to timeout when add observer into suspect ips, server: {}, tryTimes: {}",
addr.toString(), retryTimes, e);
} catch (InterruptedException e) {
// do not throw exception to user layer
// next background task will continue to remove it
logger.warn("waiting to get lock while interrupted by other threads", e);
}
}
logger.debug("removed server from suspect list: {}", addr);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
logger.debug("interrupted while removing observer from suspect ips, server: {}", addr);
} finally {
if (acquired) {
lock.unlock();
Expand All @@ -294,6 +347,10 @@ private void removeFromSuspectIPs(ObServerAddr addr) {
private void calcFailureOrClearCache(ObServerAddr addr) {
TableRoute tableRoute = tableClient.getTableRoute();
SuspectObServer server = suspectServers.get(addr);
if (server == null) {
logger.debug("skip failure calculation for removed suspect server: {}", addr);
return;
}
server.incrementFailure();
int failure = server.getFailure();
if (failure >= failureLimit) {
Expand All @@ -304,6 +361,18 @@ private void calcFailureOrClearCache(ObServerAddr addr) {
addr, failure);
}

int suspectServerCount() {
return suspectServers.size();
}

boolean containsSuspectServer(ObServerAddr addr) {
return suspectServers.containsKey(addr);
}

int lastAccessTimestampCount() {
return serverLastAccessTimestamps.size();
}

public static class SuspectObServer {
private final ObServerAddr addr;
private final long accessTimestamp;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Consumer;

import com.alipay.oceanbase.rpc.ObTableClient;
import com.alipay.oceanbase.rpc.exception.ObTableCloseException;
Expand All @@ -36,6 +37,7 @@ public class TableRoster {
private Properties properties = new Properties();
private Map<String, Object> tableConfigs = new HashMap<>();
private ObTableClientType clientType;
private Consumer<ObServerAddr> failureHandler;
/*
* ServerAddr(all) -> ObTableConnection
*/
Expand Down Expand Up @@ -65,6 +67,9 @@ public void setProperties(Properties properties) {
public void setTableConfigs(Map<String, Object> tableConfigs) {
this.tableConfigs = tableConfigs;
}
public void setFailureHandler(Consumer<ObServerAddr> failureHandler) {
this.failureHandler = failureHandler;
}
public ObTable getTable(ObServerAddr addr) {
return tables.get(addr);
}
Expand Down Expand Up @@ -101,7 +106,8 @@ public List<ObServerAddr> refreshTablesAndGetNewServers(List<ReplicaLocation> ne

ObTable obTable = new ObTable.Builder(addr.getIp(), addr.getSvrPort()) //
.setLoginInfo(tenantName, userName, password, database, clientType) //
.setProperties(properties).setConfigs(tableConfigs).setObServerAddr(addr).build();
.setProperties(properties).setConfigs(tableConfigs).setObServerAddr(addr)
.setFailureHandler(failureHandler).build();
ObTable oldObTable = tables.putIfAbsent(addr, obTable);
logger.warn("add new table addr, {}", addr.toString());
if (oldObTable != null) { // maybe create two ob table concurrently, close current ob table
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ public class TableRoute {
private IndexLocations indexLocations = null; // global index location
private TableGroupCache tableGroupCache = null;
private OdpInfo odpInfo = null;
private RouteTableRefresher routeRefresher = null;
private volatile RouteTableRefresher routeRefresher = null;
private long lastRefreshMetadataTimestamp = -1;
public final Lock refreshTableRosterLock = new ReentrantLock();

Expand Down Expand Up @@ -348,7 +348,8 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized,
tableClient.getPassword(), tableClient.getDatabase(),
tableClient.getClientType(runningMode))
.setProperties(tableClient.getProperties())
.setConfigs(tableClient.getTableConfigs()).setObServerAddr(addr).build();
.setConfigs(tableClient.getTableConfigs()).setObServerAddr(addr)
.setFailureHandler(this::reportObServerFailure).build();
addr2Table.put(addr, obTable);
servers.add(addr);
} catch (Exception e) {
Expand All @@ -363,6 +364,7 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized,
tableClient.getUserName(), tableClient.getPassword(), tableClient.getDatabase(),
tableClient.getClientType(runningMode), tableClient.getProperties(),
tableClient.getTableConfigs());
this.tableRoster.setFailureHandler(this::reportObServerFailure);
this.tableRoster.setTables(addr2Table);
this.serverRoster.reset(servers);

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

public void launchRouteRefresher() {
routeRefresher = new RouteTableRefresher(tableClient, sysUA);
routeRefresher.refreshActiveServers(serverRoster.getMembers());
routeRefresher.start();
}

public void reportObServerFailure(ObServerAddr addr) {
RouteTableRefresher refresher = routeRefresher;
if (refresher != null) {
refresher.addIntoSuspectIPs(addr);
}
}

public void removeObServer(ObServerAddr addr) {
logger.debug("remove useless table addr, {}", addr.toString());
ConcurrentHashMap<ObServerAddr, ObTable> tables = this.tableRoster.getTables();
Expand Down Expand Up @@ -480,6 +490,10 @@ public void refreshRosterByRsList(List<ObServerAddr> newRsList) throws Exception
// update new ob table and get new server address
List<ObServerAddr> servers = tableRoster.refreshTablesAndGetNewServers(replicaLocations);
serverRoster.reset(servers);
RouteTableRefresher refresher = routeRefresher;
if (refresher != null) {
refresher.refreshActiveServers(servers);
}

// 2. Get Server LDC info for weak read consistency.
success = false;
Expand Down
Loading
Loading