在上篇博客中,寫了一集群限流的Demo,這篇來分析記錄一下集群限流的原理。
不管是集群Client,或者是Server,都會(huì)實(shí)現(xiàn)TokenService服務(wù),Server端如果是內(nèi)嵌TokenService服務(wù),則默認(rèn)使用DefaultEmbeddedTokenServer,而Client端則會(huì)使用DefaultClusterTokenClient,類圖如下:

diagram.png
import java.util.Collection;
public interface TokenService {
//獲取令牌Token, 參數(shù)規(guī)則Id,獲取令牌數(shù),優(yōu)先級(jí)
TokenResult requestToken(Long ruleId, int acquireCount, boolean prioritized);
TokenResult requestParamToken(Long ruleId, int acquireCount, Collection<Object> params);
}
在服務(wù)端獲取令牌的時(shí)候,實(shí)質(zhì)是通過 DefaultEmbeddedTokenServer#requestToken獲取Token
public class DefaultEmbeddedTokenServer implements EmbeddedClusterTokenServer {
private final TokenService tokenService = TokenServiceProvider.getService();
private final ClusterTokenServer server = new SentinelDefaultTokenServer(true);
@Override
public TokenResult requestToken(Long ruleId, int acquireCount, boolean prioritized) {
if (tokenService != null) {
return tokenService.requestToken(ruleId, acquireCount, prioritized);
}
return new TokenResult(TokenResultStatus.FAIL);
}
@Override
public TokenResult requestParamToken(Long ruleId, int acquireCount, Collection<Object> params) {
if (tokenService != null) {
return tokenService.requestParamToken(ruleId, acquireCount, params);
}
return new TokenResult(TokenResultStatus.FAIL);
}
}
public class DefaultTokenService implements TokenService {
//獲取令牌
@Override
public TokenResult requestToken(Long ruleId, int acquireCount, boolean prioritized) {
//判斷是否是有效的請(qǐng)求
if (notValidRequest(ruleId, acquireCount)) {
return badRequest();
}
// 根據(jù)RuleId查詢FlowRule
FlowRule rule = ClusterFlowRuleManager.getFlowRuleById(ruleId);
if (rule == null) {
return new TokenResult(TokenResultStatus.NO_RULE_EXISTS);
}
//獲取令牌
return ClusterFlowChecker.acquireClusterToken(rule, acquireCount, prioritized);
}
//判斷是否是一個(gè)有效的請(qǐng)求
private boolean notValidRequest(Long id, int count) {
return id == null || id <= 0 || count <= 0;
}
}
//ClusterFlowChecker.java
static TokenResult acquireClusterToken(/*@Valid*/ FlowRule rule, int acquireCount, boolean prioritized) {
Long id = rule.getClusterConfig().getFlowId();
//是否繼續(xù),根據(jù)RuleId,獲取NameSpace,根據(jù)nameSpace,判斷nameSpace限流是否通過
if (!allowProceed(id)) {
return new TokenResult(TokenResultStatus.TOO_MANY_REQUEST);
}
ClusterMetric metric = ClusterMetricStatistics.getMetric(id);
if (metric == null) {
return new TokenResult(TokenResultStatus.FAIL);
}
//獲取Metric,滑動(dòng)窗口實(shí)現(xiàn),這里獲取的是通過的請(qǐng)求平均值
double latestQps = metric.getAvg(ClusterFlowEvent.PASS_REQUEST);
//獲取全局閥值 根據(jù)規(guī)則判斷是否為全局限流還是平均分?jǐn)?,并獲取總的閥值
double globalThreshold = calcGlobalThreshold(rule) * ClusterServerConfigManager.getExceedCount();
//判斷剩余請(qǐng)求數(shù)
double nextRemaining = globalThreshold - latestQps - acquireCount;
//如果>=0,則代表請(qǐng)求可以通過
if (nextRemaining >= 0) {
//記錄請(qǐng)求數(shù)量
metric.add(ClusterFlowEvent.PASS, acquireCount);
metric.add(ClusterFlowEvent.PASS_REQUEST, 1);
if (prioritized) {
metric.add(ClusterFlowEvent.OCCUPIED_PASS, acquireCount);
}
return new TokenResult(TokenResultStatus.OK)
.setRemaining((int) nextRemaining)
.setWaitInMs(0);
} else {
//這里忽略優(yōu)先級(jí)邏輯
//其他直接返回失敗
metric.add(ClusterFlowEvent.BLOCK, acquireCount);
metric.add(ClusterFlowEvent.BLOCK_REQUEST, 1);
ClusterServerStatLogUtil.log("flow|block|" + id, acquireCount);
ClusterServerStatLogUtil.log("flow|block_request|" + id, 1);
if (prioritized) {
metric.add(ClusterFlowEvent.OCCUPIED_BLOCK, acquireCount);
ClusterServerStatLogUtil.log("flow|occupied_block|" + id, 1);
}
return blockedResult();
}
}
static boolean allowProceed(long flowId) {
String namespace = ClusterFlowRuleManager.getNamespace(flowId);
return GlobalRequestLimiter.tryPass(namespace);
}
static boolean allowProceed(long flowId) {
String namespace = ClusterFlowRuleManager.getNamespace(flowId);
return GlobalRequestLimiter.tryPass(namespace);
}
public static boolean tryPass(String namespace) {
if (namespace == null) {
return false;
}
RequestLimiter limiter = GLOBAL_QPS_LIMITER_MAP.get(namespace);
if (limiter == null) {
return true;
}
return limiter.tryPass();
}
public boolean tryPass() {
if (canPass()) {
add(1);
return true;
}
return false;
}
ClusterServerConfigManager.loadGlobalFlowConfig配置了nameSpace對(duì)應(yīng)的ServerFlowConfig
而在客戶端的時(shí)候,通過netty通信發(fā)送到服務(wù)端,由服務(wù)端驗(yàn)證是否通過。
@Override
public TokenResult requestToken(Long flowId, int acquireCount, boolean prioritized) {
//驗(yàn)證是否有效請(qǐng)求
if (notValidRequest(flowId, acquireCount)) {
return badRequest();
}
//初始化FlowRequest
FlowRequestData data = new FlowRequestData().setCount(acquireCount)
.setFlowId(flowId).setPriority(prioritized);
ClusterRequest<FlowRequestData> request = new ClusterRequest<>(ClusterConstants.MSG_TYPE_FLOW, data);
try {
//發(fā)送請(qǐng)求到服務(wù)端
TokenResult result = sendTokenRequest(request);
logForResult(result);
return result;
} catch (Exception ex) {
ClusterClientStatLogUtil.log(ex.getMessage());
return new TokenResult(TokenResultStatus.FAIL);
}
}
private TokenResult sendTokenRequest(ClusterRequest request) throws Exception {
if (transportClient == null) {
RecordLog.warn(
"[DefaultClusterTokenClient] Client not created, please check your config for cluster client");
return clientFail();
}
ClusterResponse response = transportClient.sendRequest(request);
TokenResult result = new TokenResult(response.getStatus());
if (response.getData() != null) {
FlowTokenResponseData responseData = (FlowTokenResponseData)response.getData();
result.setRemaining(responseData.getRemainingCount())
.setWaitInMs(responseData.getWaitInMs());
}
return result;
}
在FlowSlot限流的時(shí)候,根據(jù)節(jié)點(diǎn)配置是否啟用ClusterMode,判斷是否走限流,然后根據(jù)節(jié)點(diǎn)狀態(tài)(是Server,或者是Client)獲取服務(wù),申請(qǐng)令牌。
static boolean passCheck(/*@NonNull*/ FlowRule rule, Context context, DefaultNode node, int acquireCount,
boolean prioritized) {
String limitApp = rule.getLimitApp();
if (limitApp == null) {
return true;
}
//如果是集群模式
if (rule.isClusterMode()) {
return passClusterCheck(rule, context, node, acquireCount, prioritized);
}
//單機(jī)模式
return passLocalCheck(rule, context, node, acquireCount, prioritized);
}
private static boolean passClusterCheck(FlowRule rule, Context context, DefaultNode node, int acquireCount,
boolean prioritized) {
try {
//獲取令牌服務(wù)
TokenService clusterService = pickClusterService();
if (clusterService == null) {
return fallbackToLocalOrPass(rule, context, node, acquireCount, prioritized);
}
long flowId = rule.getClusterConfig().getFlowId();
//申請(qǐng)令牌
TokenResult result = clusterService.requestToken(flowId, acquireCount, prioritized);
return applyTokenResult(result, rule, context, node, acquireCount, prioritized);
} catch (Throwable ex) {
RecordLog.warn("[FlowRuleChecker] Request cluster token unexpected failed", ex);
}
//如果失敗,降級(jí)為單機(jī)模式
return fallbackToLocalOrPass(rule, context, node, acquireCount, prioritized);
}
private static TokenService pickClusterService() {
if (ClusterStateManager.isClient()) {
return TokenClientProvider.getClient();
}
if (ClusterStateManager.isServer()) {
return EmbeddedClusterTokenServerProvider.getServer();
}
return null;
}