Sentinel 源码精读
概述
Sentinel 的核心是 ProcessorSlotChain(责任链模式),每个 Slot 负责一种规则检查。理解这个链路是掌握 Sentinel 的关键。
核心执行流程
@SentinelResource
↓
SentinelResourceAspect.invoke()
↓
SphU.entry(resource, entryType)
↓
CtSph.entryWithPriority()
↓
CtSph.entryForFlow() / entryForDegrade()
↓
ProcessorSlotChain.entry()
↓
责任链逐级执行
├── NodeSelectorSlot → 构建调用树
├── ClusterBuilderSlot → 集群节点统计
├── LogSlot → 日志记录
├── StatisticSlot → 统计指标(滑动窗口)
│ ├── FlowSlot → 流控检查
│ ├── DegradeSlot → 熔断检查
│ ├── AuthoritySlot → 黑白名单
│ ├── SystemSlot → 系统保护
│ └── ParamFlowSlot → 热点参数限流
└── 检查通过 → 执行业务代码一、SentinelResourceAspect
1.1 切入点
java
// com.alibaba.csp.sentinel.annotation.aspectj.SentinelResourceAspect
@Aspect
public class SentinelResourceAspect {
// 拦截所有 @SentinelResource 注解的方法
@Around("@annotation(sentinelResource)")
public Object invokeResource(ProceedingJoinPoint pjp,
SentinelResource sentinelResource) throws Throwable {
// 1. 解析注解
String resourceName = getResourceName(sentinelResource.value(), pjp);
EntryType entryType = sentinelResource.entryType();
int entryPriority = sentinelResource.priority();
// 2. 构建资源入口
try (Entry entry = SphU.entry(resourceName, entryType)) {
// 3. 执行业务方法
return pjp.proceed();
} catch (BlockException ex) {
// 4. 被限流/降级 → 调用 blockHandler
return handleBlockException(pjp, sentinelResource, ex);
} catch (Throwable ex) {
// 5. 业务异常 → 调用 fallback
return handleFallback(pjp, sentinelResource, ex);
}
}
}1.2 blockHandler 查找逻辑
java
private Object handleBlockException(ProceedingJoinPoint pjp,
SentinelResource annotation,
BlockException ex) throws Throwable {
// 优先找 blockHandlerClass 中的 blockHandler 方法
if (annotation.blockHandlerClass() != Void.class
&& !annotation.blockHandler().isEmpty()) {
Class<?>[] argTypes = getArgTypes(pjp, BlockException.class);
Method method = annotation.blockHandlerClass()
.getDeclaredMethod(annotation.blockHandler(), argTypes);
if (!Modifier.isStatic(method.getModifiers())) {
throw new IllegalArgumentException("blockHandler must be static");
}
return method.invoke(null, getArgs(pjp, ex));
}
// 否则找当前类中的 blockHandler 方法
// ...
throw ex; // 未配置 blockHandler,抛出 BlockException
}二、CtSph 核心
2.1 资源入口
java
// com.alibaba.csp.sentinel.CtSph
public class CtSph implements Sph {
@Override
public Entry entry(String name, EntryType type) throws BlockException {
return entryWithPriority(name, type, false);
}
private Entry entryWithPriority(String name, EntryType type,
boolean prioritize) throws BlockException {
// 1. 从缓存获取 ProcessorSlotChain
ProcessorSlot<Object> chain = lookProcessChain(name);
// 2. 创建 Entry
CtEntry entry = new CtEntry(name, type);
try {
// 3. 执行 SlotChain(核心!)
chain.entry(entry.getContext(), entry, name, type);
} catch (BlockException e) {
// 4. 被限流/降级
entry.exit();
throw e;
}
return entry;
}
}2.2 lookProcessChain
java
private ProcessorSlot<Object> lookProcessChain(String resourceName) {
// 双重检查锁缓存
ProcessorSlot<Object> chain = chainMap.get(resourceName);
if (chain == null) {
synchronized (LOCK) {
chain = chainMap.get(resourceName);
if (chain == null) {
// 构建 SlotChain(每个资源独立构建)
chain = SlotChainProvider.newSlotChain();
chainMap.put(resourceName, chain);
}
}
}
return chain;
}三、ProcessorSlotChain 构建
3.1 SlotChainProvider
java
// com.alibaba.csp.sentinel.slotchain.SlotChainProvider
public final class SlotChainProvider {
public static ProcessorSlotChain newSlotChain() {
// SPI 加载 SlotChainBuilder
SlotChainBuilder builder = SpiLoader
.loadFirstInstanceOrDefault(
SlotChainBuilder.class,
DefaultSlotChainBuilder.class
);
return builder.build();
}
}3.2 DefaultSlotChainBuilder
java
public class DefaultSlotChainBuilder implements SlotChainBuilder {
@Override
public ProcessorSlotChain build() {
ProcessorSlotChain chain = new DefaultProcessorSlotChain();
// SPI 加载所有 ProcessorSlot(责任链节点)
List<ProcessorSlot> slots = SpiLoader
.loadInstanceListSorted(ProcessorSlot.class);
// 默认加载顺序:
// NodeSelectorSlot → ClusterBuilderSlot → LogSlot
// → StatisticSlot → FlowSlot → DegradeSlot
// → AuthoritySlot → SystemSlot → ParamFlowSlot
for (ProcessorSlot slot : slots) {
chain.addLast(slot);
}
return chain;
}
}四、StatisticSlot 滑动窗口
4.1 统计数据
java
// com.alibaba.csp.sentinel.slots.statistic.StatisticSlot
public class StatisticSlot extends AbstractLinkedProcessorSlot<DefaultNode> {
@Override
public void entry(Context context, ResourceWrapper resourceWrapper,
DefaultNode node, int count, boolean prioritized, Object... args)
throws Throwable {
try {
// 1. 执行后续 Slot(FlowSlot / DegradeSlot 等)
fireEntry(context, resourceWrapper, node, count, prioritized, args);
// 2. 请求通过 → 增加通过计数
node.increaseThreadNum();
node.addPassRequest(count);
// 3. 集群节点统计
if (resourceWrapper.getEntryType() != EntryType.IN) {
Constants.ENTRY_NODE.increaseThreadNum();
Constants.ENTRY_NODE.addPassRequest(count);
}
} catch (BlockException e) {
// 3. 请求被限流 → 增加阻塞计数
node.increaseBlockedQps(count);
throw e;
} catch (Throwable e) {
// 4. 业务异常 → 增加异常计数
node.increaseExceptionQps(count);
throw e;
}
}
}4.2 LeapArray 滑动窗口
java
// com.alibaba.csp.sentinel.util.LeapArray
// 核心数据结构:时间片轮转 + 滑动平均
public class LeapArray<T> {
// 窗口数组
protected final AtomicReferenceArray<WindowWrap<T>> array;
// 每个时间片的长度(毫秒),默认 500ms
protected int windowLengthMs;
// 时间片数量,默认 20(总共 10s 滑动窗口)
protected int sampleCount;
// 总滑动窗口时间,默认 10000ms
protected int intervalInMs;
// 获取当前时间所属的窗口
public WindowWrap<T> currentWindow(long timeMillis) {
// 1. 计算时间片索引
int idx = calculateTimeIdx(timeMillis);
// 2. 计算窗口起始时间
long windowStart = calculateWindowStart(timeMillis);
while (true) {
WindowWrap<T> old = array.get(idx);
if (old == null) {
// 3.1 该槽位为空 → 新建窗口
WindowWrap<T> window = new WindowWrap<>(
windowLengthMs, windowStart, newEmptyBucket(timeMillis));
if (array.compareAndSet(idx, null, window)) {
return window;
}
Thread.yield();
} else if (windowStart == old.windowStart()) {
// 3.2 命中当前窗口 → 直接返回
return old;
} else if (windowStart > old.windowStart()) {
// 3.3 窗口已过期 → CAS 重置
if (array.compareAndSet(idx, old,
new WindowWrap<>(windowLengthMs, windowStart, newEmptyBucket(timeMillis)))) {
return array.get(idx);
}
Thread.yield();
} else {
// 不会发生
return old;
}
}
}
}4.3 滑动窗口统计
java
// 计算当前时间窗口内的总请求数
public long totalCount() {
long count = 0;
long currentTime = TimeUtil.currentTimeMillis();
for (int i = 0; i < array.length(); i++) {
WindowWrap<MetricBucket> wrap = array.get(i);
if (wrap != null && !wrap.isDeprecated(currentTime)) {
// 只累加未过期的窗口
count += wrap.value().pass();
}
}
return count;
}五、FlowSlot 流控检查
5.1 核心入口
java
// com.alibaba.csp.sentinel.slots.block.flow.FlowSlot
public class FlowSlot extends AbstractLinkedProcessorSlot<DefaultNode> {
@Override
public void entry(Context context, ResourceWrapper resourceWrapper,
DefaultNode node, int count, boolean prioritized, Object... args)
throws Throwable {
// 1. 检查流控规则
checkFlow(resourceWrapper, context, node, count);
// 2. 通过 → 执行下一个 Slot
fireEntry(context, resourceWrapper, node, count, prioritized, args);
}
void checkFlow(ResourceWrapper resource, Context context,
DefaultNode node, int count) throws BlockException {
// 获取该资源的所有流控规则
List<FlowRule> rules = FlowRuleManager.getRules(resource.getName());
if (rules == null) return;
for (FlowRule rule : rules) {
// 逐一检查每条规则
if (!canPassCheck(rule, context, node, count, false)) {
throw new FlowException(rule.getLimitApp(), rule);
}
}
}
}5.2 FlowRuleChecker
java
// com.alibaba.csp.sentinel.slots.block.flow.FlowRuleChecker
public class FlowRuleChecker {
public boolean canPassCheck(FlowRule rule, Context context,
DefaultNode node, int count, boolean prioritized) {
// 1. 获取限流统计节点
DefaultNode selectedNode = selectNodeByStrategy(rule, context, node);
return passLocalCheck(rule, selectedNode, count, prioritized);
}
private boolean passLocalCheck(FlowRule rule, DefaultNode node,
int count, boolean prioritized) {
// 2. 从规则获取限流检查器
TrafficShapingController controller = rule.getRater();
// 3. 执行限流逻辑
return controller.canPass(node, count, prioritized);
}
}5.3 限流算法实现
java
// DefaultController(快速失败)
public class DefaultController implements TrafficShapingController {
@Override
public boolean canPass(Node node, int acquireCount, boolean prioritized) {
// 当前窗口已通过数 + 本次请求数 <= 阈值
return node.passQps() + acquireCount <= rule.getCount();
}
}
// WarmUpController(冷启动)
public class WarmUpController implements TrafficShapingController {
private long storedTokens; // 存储的令牌数
private long maxTokenPerSlice; // 每秒最大令牌数
@Override
public boolean canPass(Node node, int acquireCount, boolean prioritized) {
// 冷启动阶段 token 逐步增加
long currentQps = node.passQps();
long expectQps = rule.getCount();
// 计算当前斜率
long tokens = Math.min(storedTokens + currentQps, maxTokenPerSlice);
if (tokens < maxTokenPerSlice) {
storedTokens = tokens;
return true;
}
return false;
}
}
// RateLimiterController(排队等待)
public class RateLimiterController implements TrafficShapingController {
private final AtomicLong lastPassTime = new AtomicLong(0);
@Override
public boolean canPass(Node node, int acquireCount, boolean prioritized) {
long currentTime = TimeUtil.currentTimeMillis();
long costTime = Math.round(1.0 * 1000 / rule.getCount()); // 每个请求间的时间间隔
long expectedTime = costTime + lastPassTime.get();
if (expectedTime <= currentTime) {
// 可以立即通过
lastPassTime.set(currentTime);
return true;
} else {
// 需要等待
long waitTime = costTime + lastPassTime.get() - currentTime;
if (waitTime > rule.getMaxQueueingTimeMs()) {
return false; // 超时,拒绝
}
lastPassTime.addAndGet(costTime);
return true;
}
}
}六、DegradeSlot 熔断检查
6.1 熔断逻辑
java
// com.alibaba.csp.sentinel.slots.block.degrade.DegradeSlot
public class DegradeSlot extends AbstractLinkedProcessorSlot<DefaultNode> {
@Override
public void entry(Context context, ResourceWrapper resourceWrapper,
DefaultNode node, int count, boolean prioritized, Object... args)
throws Throwable {
// 检查熔断规则
DegradeRuleManager.checkDegrade(resourceWrapper, context, node, count);
fireEntry(context, resourceWrapper, node, count, prioritized, args);
}
}6.2 CircuitBreaker 状态机
java
// com.alibaba.csp.sentinel.slots.block.degrade.circuitbreaker.CircuitBreaker
public abstract class CircuitBreaker {
// 三种状态
private volatile double minRequestAmount; // 最小请求数
private volatile double statIntervalMs; // 统计时间窗口
private volatile double timeWindow; // 熔断持续时间
public abstract void onRequestComplete(Context context);
public abstract boolean tryPass(Context context);
}
// ExceptionCircuitBreaker(异常比例/异常数)
public class ExceptionCircuitBreaker extends CircuitBreaker {
@Override
public boolean tryPass(Context context) {
// OPEN → 快速失败
if (currentState.get() == State.OPEN) {
// 检查是否达到半开时间
if (System.currentTimeMillis() >= nextRetryTimestamp) {
// OPEN → HALF_OPEN
currentState.compareAndSet(State.OPEN, State.HALF_OPEN);
return true;
}
return false; // 熔断中,拒绝请求
}
// CLOSED / HALF_OPEN → 允许通过
return true;
}
@Override
public void onRequestComplete(Context context) {
// 每个请求完成后更新统计
Entry entry = context.getCurEntry();
// 统计:总请求数、异常数、异常比例
// 判断是否触发熔断
if (entry.getErrorCount() > 0) {
// 检查异常比例是否超过阈值
double errorRatio = entry.getErrorCount() / entry.getTotalCount();
if (errorRatio > rule.getCount()) {
// CLOSED → OPEN
currentState.set(State.OPEN);
nextRetryTimestamp = System.currentTimeMillis()
+ (long) rule.getTimeWindow() * 1000;
}
}
}
}七、规则持久化与数据源
规则在内存中加载后重启即丢失,生产环境必须持久化。Sentinel 通过 DataSource SPI 支持从外部存储加载规则并动态推送。
DataSource 抽象
外部存储(Nacos/Redis/DB/文件)
│ DataSource(读 + 注册监听)
▼
Property(规则属性,观察者模式)
│ 变更通知
▼
RuleManager(FlowRuleManager / DegradeRuleManager ...)
│
▼
责任链中的 Slot 读取最新规则java
// 数据源接口
public interface DataSource<T, S> {
T loadConfig() throws Exception; // 初始加载
void close() throws Exception;
void addProperty(Property<S> property); // 注册规则属性监听
}推模式(推荐):Nacos 数据源
Nacos 配置变更 → DataSource 监听 → 规则刷新(实时推送)java
public class NacosDataSourceInit {
public void init() {
// 1. 连接 Nacos 配置中心
ConfigService configService = NacosFactory.createConfigService(
"127.0.0.1:8848");
// 2. 创建数据源:监听 dataId = "flow-rules" 的配置
String groupId = "SENTINEL_GROUP";
DataSource<String, List<FlowRule>> ds = new NacosDataSource<>(
configService, groupId, "flow-rules",
source -> JSON.parseArray(source, FlowRule.class));
// 3. 注册到规则管理器(变更自动生效)
FlowRuleManager.register2Property(ds.getProperty());
}
}拉模式 vs 推模式
| 模式 | 原理 | 实时性 | 一致性 |
|---|---|---|---|
| 拉模式 | 定时从存储拉取(如 file/redis) | 秒级延迟 | 多节点可能不一致 |
| 推模式 | 存储变更主动推送(Nacos/APOLLO) | 实时 | 所有节点一致 |
规则在责任链中的读取
java
// FlowSlot 读取规则的源码路径
public class FlowSlot extends AbstractLinkedProcessorSlot<DefaultNode> {
@Override
public void entry(Context context, ResourceWrapper resourceWrapper,
DefaultNode node, int count, boolean prioritized) {
// 从 FlowRuleManager 读取当前资源的最新规则
// (规则由数据源动态刷新,这里总是拿到最新)
checkFlow(resourceWrapper, context, node, count, prioritized);
}
}八、总结
| 组件 | 职责 | 关键类 |
|---|---|---|
| 切面 | 拦截 @SentinelResource | SentinelResourceAspect |
| 入口 | 创建 Entry + 执行 SlotChain | CtSph |
| 责任链 | SPI 加载所有 Slot | ProcessorSlotChain、DefaultSlotChainBuilder |
| 统计 | 滑动窗口 + QPS 计数 | StatisticSlot、LeapArray、MetricBucket |
| 流控 | 快速失败 / WarmUp / 排队 | FlowSlot、FlowRuleChecker、TrafficShapingController |
| 熔断 | 异常比例 / RT / 异常数 | DegradeSlot、ExceptionCircuitBreaker、ResponseTimeCircuitBreaker |
| 黑名单 | 来源鉴权 | AuthoritySlot |
| 系统保护 | Load / CPU / 线程 | SystemSlot |
| 热点 | 参数级别限流 | ParamFlowSlot |
参考链接: