Dubbo 源码阅读
SPI(扩展点加载)
核心机制
Dubbo 的 SPI 是 JDK SPI 的增强版,支持按需加载、依赖注入、自适应扩展。
java
// 接口定义 @SPI
@SPI("dubbo")
public interface Protocol {
<T> Exporter<T> export(Invoker<T> invoker);
<T> Invoker<T> refer(Class<T> type, URL url);
}配置文件
META-INF/dubbo/org.apache.dubbo.rpc.Protocol
dubbo=org.apache.dubbo.rpc.protocol.dubbo.DubboProtocol
rest=org.apache.dubbo.rpc.protocol.rest.RestProtocol
tri=org.apache.dubbo.rpc.protocol.tri.TripleProtocol自适应扩展
java
@Adaptive
public class AdaptiveExtensionFactory implements ExtensionFactory {
// 运行时根据 URL 参数决定使用哪个实现
@Override
public <T> T getExtension(Class<T> type, String name) {
ExtensionFactory factory = getAdaptiveExtension();
return factory.getExtension(type, name);
}
}服务暴露
流程
ServiceBean (Spring 容器初始化)
│
└─ ServiceConfig.export()
│
├─ 检查配置(接口、注册中心、协议)
│
├─ ProxyFactory.getInvoker()
│ └─ 将 Bean 封装成 Invoker
│
├─ Protocol.export(invoker)
│ │
│ ├─ DubboProtocol.openServer()
│ │ └─ 启动 Netty Server(默认 20880)
│ │
│ └─ RegistryProtocol.export()
│ │
│ ├─ 将 URL 注册到注册中心
│ └─ subscribeOverride() ← 订阅配置变更
│
└─ 暴露完成,等待调用Invoker 核心
java
public class DubboInvoker<T> extends AbstractInvoker<T> {
private final ExchangeClient client;
@Override
protected Result doInvoke(Invocation invocation) throws Throwable {
// 1. 构建 RPC 请求
RpcInvocation rpcInv = new RpcInvocation(invocation);
// 2. 选择客户端(连接复用)
ExchangeClient currentClient = clients.get(ThreadLocalRandom.current().nextInt(clients.size()));
// 3. 发送请求
CompletableFuture<AppResponse> future =
currentClient.request(rpcInv);
// 4. 等待响应(同步/异步)
return future.get(invocation.getTimeout(), TimeUnit.MILLISECONDS);
}
}服务引用
流程
@DubboReference private UserService userService;
│
└─ ReferenceConfig.get()
│
├─ 从注册中心获取服务 URL 列表
│
├─ 通过 Router 路由过滤
│
├─ 通过 LoadBalance 选择一个 Provider
│
├─ Protocol.refer() → 创建 Invoker
│ └─ DubboProtocol.openClient()
│ └─ 建立 Netty 连接
│
└─ ProxyFactory.getProxy(invoker)
└─ 生成代理对象(Javassist / JDK)集群容错
容错策略
| 策略 | 说明 | 适用场景 |
|---|---|---|
| Failover | 失败重试其他节点(默认) | 幂等操作 |
| Failfast | 立即失败 | 非幂等操作 |
| Failsafe | 忽略失败 | 日志上报 |
| Failback | 失败后定时重试 | 消息通知 |
| Forking | 并行调用多个,取最快返回 | 实时性要求高 |
| Broadcast | 广播调用所有 | 状态同步 |
源码实现
java
public class FailoverClusterInvoker<T> extends AbstractClusterInvoker<T> {
@Override
public Result doInvoke(Invocation invocation, List<Invoker<T>> invokers, LoadBalance loadbalance) {
int retries = getUrl().getParameter("retries", 2); // 重试 2 次
for (int i = 0; i <= retries; i++) {
Invoker<T> invoker = select(loadbalance, invocation, invokers, selected);
try {
return invoker.invoke(invocation);
} catch (RpcException e) {
if (i >= retries) throw e; // 重试耗尽
}
}
throw new RpcException("所有节点都失败了");
}
}