跳到主要内容

多线程和Reactor模式下解决数据丢失

到了此章节,相信你对 link-flow 的整体流程有了大致的掌握了,但在整个流程中,对于线程方面有很多的细节,稍微处理不好就容易丢失,所以此章节将会讲清楚这些细节问题。

多线程的切换

首先来看这四个部分,在这四个部分中都进行了线程名的输出

第1处:服务过滤执行的开始

public Mono<Response<ServiceInstance>> choose(Request request) {
ServiceInstanceListSupplier supplier = serviceInstanceListSupplierProvider
.getIfAvailable(NoopServiceInstanceListSupplier::new);
System.out.println("第一处执行的线程;"+Thread.currentThread().getName());
return supplier.get(request).next()
.map(serviceInstances -> processInstanceResponse(supplier, serviceInstances));
}

第2处:获取到所有服务列表 第3处:服务的路由过滤

/**
* 当要调用服务的时候,会调用此方法,此方法会返回所有的服务列表。路由过滤器会在此处执行
* */
@Override
public Flux<List<ServiceInstance>> get() {
//获取所有的服务列表
Flux<List<ServiceInstance>> listFlux = super.get();
//从ThreadLocal获取参数
Map<String, Object> parameterMap = BaseParameterHolder.getParameterMap();
Map<String,Object> newMap = new HashMap<>(BaseParameterHolder.getParameterMap().size());
newMap.putAll(parameterMap);
System.out.println("第2处执行的线程;"+Thread.currentThread().getName());
listFlux = listFlux.map(serviceInstances -> {
System.out.println("第3处执行的线程;"+Thread.currentThread().getName());
//到这里线程已经发生变化了,所以要把之前线程的将参数放入到ThreadLocal中,这样才能在后续的过滤器中获取到参数
BaseParameterHolder.setParameterMap(newMap);
List<ServiceInstance> allServers = new ArrayList<>();
Optional.ofNullable(serviceInstances).ifPresent(allServers::addAll);
//执行过滤器
linkFlowFilterLoadBalance.selectServer(allServers);
return allServers;
});
//返回结果
return listFlux;
}

第4处:服务版本权重的执行

private Response<ServiceInstance> getInstanceResponse(List<ServiceInstance> instances) {
if (instances.isEmpty()) {
if (log.isWarnEnabled()) {
log.warn("No servers available for service: " + serviceId);
}
return new EmptyResponse();
}

System.out.println("第四处执行的线程;"+Thread.currentThread().getName());

//经过权重过滤选择后,肯定是一个服务实例了
WeightInfoWrapper weightInfoWrapper = linkFlowWeight.parseWeightInfo();
if (linkFlowWeight.isServiceWeight(instances,weightInfoWrapper)) {
ServiceInstance serviceInstance = linkFlowWeight.selectServiceInstance(instances, weightInfoWrapper);
instances.clear();
instances.add(serviceInstance);
}


// Do not move position when there is only 1 instance, especially some suppliers
// have already filtered instances
if (instances.size() == 1) {
return new DefaultResponse(instances.get(0));
}

// Ignore the sign bit, this allows pos to loop sequentially from 0 to
// Integer.MAX_VALUE
int pos = this.position.incrementAndGet() & Integer.MAX_VALUE;

ServiceInstance instance = instances.get(pos % instances.size());

return new DefaultResponse(instance);
}

付费内容提示

该文档的全部内容仅对「码力全开」项目实战&技术讲解 知识星球用户开放

加入星球,一次获得完整项目资料、全栈技术知识库和长期答疑服务。

100万+字全栈技术知识库深入讲解技术核心、数据库、中间件和分布式等内容
8套热门的实战项目持续更新的企业级项目覆盖高并发、微服务、数据中台 和 AI Agent 等方向
AI 技术知识大模型面试详解覆盖 AI 模型原理、Agent、RAG、MCP、Skills、Harness 等核心知识点
文档 + 视频两种讲解形式既能系统阅读,也能跟随视频理解核心业务

完整项目实战资料

每套项目均包含从 0 到 1 讲解文档核心业务讲解视频

从基础项目到复杂业务场景,项目资料会持续更新。

8 套项目
  • 01Nexus Agent AI 智能体
  • 02Nexus Agent Pro 完全版
  • 03黑马点评Plus
  • 04大麦
  • 05大麦Pro
  • 06大麦AI
  • 07流量切换
  • 08数据中台

加入后还能获得

进入星球后,即可享受上述所有服务,保证不会再有其他隐藏费用。

从学习、面试到项目启动,都可以继续获得支持。

  • 1 对 1 解答项目和技术问题都可以提问
  • 针对性补充没有讲清楚的内容会继续补充
  • 面试与简历指导梳理回答技巧和项目亮点
  • 中间件云环境项目依赖可以直接接入使用
  • 面试后复盘被问住的问题可以继续交流
  • 远程问题解决项目启动问题可协助排查
知识星球二维码

扫码进入知识星球

  1. 打开微信,扫描左侧二维码,加入「码力全开」项目实战&技术讲解 知识星球
  2. 查看星球使用指导,获取完整项目讲解资料索引
解锁全部付费内容
🎁优惠