分布式简介
分布式简介
简介
分布式系统(Distributed System) 是指由多台独立计算机通过网络连接协同工作的系统,这些计算机在用户看来就像是一台单一的计算机。分布式系统的核心目标是:通过将计算、存储、网络等资源分散到多个节点上,从而突破单机的性能瓶颈,提升系统的整体处理能力、可用性和可扩展性。
分布式系统广泛应用于互联网、大数据、云计算、物联网等领域,是现代大型应用系统的基石。典型的分布式系统包括:分布式数据库(如 TiDB、OceanBase)、分布式文件系统(如 HDFS、GFS)、分布式缓存(如 Redis Cluster)、分布式消息队列(如 Kafka、RocketMQ)、分布式协调服务(如 ZooKeeper、etcd)等。
构建分布式系统面临诸多挑战:网络不可靠、时钟不同步、节点故障、数据一致性、并发控制等。理解这些挑战并掌握相应的解决方案,是设计健壮分布式系统的前提。
分布式系统的演进
罗马不是一天建成的,同理,现代分布式系统架构也不是一蹴而就的,而是逐步发展的演化过程。随着业务的不断发展,用户体量的增加,系统的复杂度势必不断攀升,最终迫使系统架构进化,以应对挑战。
了解分布式系统架构的演化过程,有利于我们了解架构进化的发展规律和业界一些成熟的应对方案。帮助我们在实际工作中,如何去思考架构,如何去凝练解决方案。
单机架构
- 场景:网站运营初期,访问用户少,一台服务器绰绰有余。
- 特征:应用程序、数据库、文件等所有的资源都在一台服务器上。
- 描述:通常服务器操作系统使用 linux,应用程序使用 PHP 开发,然后部署在 Apache 上,数据库使用 Mysql,通俗称为 LAMP。汇集各种免费开源软件以及一台廉价服务器就可以开始系统的发展之路了。

应用服务和数据服务分离
- 场景:越来越多的用户访问导致性能越来越差,越来越多的数据导致存储空间不足,一台服务器已不足以支撑。
- 特征:应用服务器、数据库服务器、文件服务器分别独立部署。
- 描述:三台服务器对性能要求各不相同:
- 应用服务器要处理大量业务逻辑,因此需要更快更强大的 CPU;
- 数据库服务器需要快速磁盘检索和数据缓存,因此需要更快的硬盘和更大的内存;
- 文件服务器需要存储大量文件,因此需要更大容量的硬盘。

使用缓存改善性能
- 场景:随着用户逐渐增多,数据库压力太大导致访问延迟。
- 特征:由于网站访问和财富分配一样遵循二八定律:80% 的业务访问集中在 20% 的数据上。将数据库中访问较集中的少部分数据缓存在内存中,可以减少数据库的访问次数,降低数据库的访问压力。
- 描述:缓存分为两种:应用服务器上的本地缓存和分布式缓存服务器上的远程缓存。
- 本地缓存访问速度更快,但缓存数据量有限,同时存在与应用程序争用内存的情况。
- 分布式缓存可以采用集群方式,理论上可以做到不受内存容量限制的缓存服务。

负载均衡
- 场景:使用缓存后,数据库访问压力得到有效缓解。但是单一应用服务器能够处理的请求连接有限,在访问高峰期,成为瓶颈。
- 特征:多台服务器通过负载均衡同时向外部提供服务,解决单一服务器处理能力和存储空间不足的问题。
- 描述:使用集群是系统解决高并发、海量数据问题的常用手段。通过向集群中追加资源,提升系统的并发处理能力,使得服务器的负载压力不再成为整个系统的瓶颈。

数据库读写分离
- 场景:网站使用缓存后,使绝大部分数据读操作访问都可以不通过数据库就能完成,但是仍有一部分读操作和全部的写操作需要访问数据库,在网站的用户达到一定规模后,数据库因为负载压力过高而成为网站的瓶颈。
- 特征:目前大部分的主流数据库都提供主从热备功能,通过配置两台数据库主从关系,可以将一台数据库服务器的数据更新同步到一台服务器上。网站利用数据库的主从热备功能,实现数据库读写分离,从而改善数据库负载压力。
- 描述:应用服务器在写操作的时候,访问主数据库,主数据库通过主从复制机制将数据更新同步到从数据库。这样当应用服务器在读操作的时候,访问从数据库获得数据。为了便于应用程序访问读写分离后的数据库,通常在应用服务器端使用专门的数据访问模块,使数据库读写分离的对应用透明。

多级缓存
- 场景:中国网络环境复杂,不同地区的用户访问网站时,速度差别也极大。
- 特征:采用 CDN 和反向代理加快系统的静态资源访问速度。
- 描述:CDN 和反向代理的基本原理都是缓存,区别在于:
- CDN 部署在网络提供商的机房,使用户在请求网站服务时,可以从距离自己最近的网络提供商机房获取数据;
- 而反向代理则部署在网站的中心机房,当用户请求到达中心机房后,首先访问的服务器时反向代理服务器,如果反向代理服务器中缓存着用户请求的资源,就将其直接返回给用户。

业务拆分
- 场景:大型网站的业务场景日益复杂,分为多个产品线。
- 特征:采用分而治之的手段将整个网站业务分成不同的产品线。系统上按照业务进行拆分改造,应用服务器按照业务区分进行分别部署。
- 描述:应用之间可以通过超链接建立关系,也可以通过消息队列进行数据分发,当然更多的还是通过访问同一个数据存储系统来构成一个关联的完整系统。
- 纵向拆分:将一个大应用拆分为多个小应用,如果新业务较为独立,那么就直接将其设计部署为一个独立的 Web 应用系统。纵向拆分相对较为简单,通过梳理业务,将较少相关的业务剥离即可。
- 横向拆分:将复用的业务拆分出来,独立部署为分布式服务,新增业务只需要调用这些分布式服务横向拆分需要识别可复用的业务,设计服务接口,规范服务依赖关系。

分库分表
- 场景:随着大型网站业务持续增长,数据库经过读写分离,从一台服务器拆分为两台服务器,依然不能满足需求。
- 特征:数据库采用分布式数据库。
- 描述:分布式数据库是数据库拆分的最后方法,只有在单表数据规模非常庞大的时候才使用。不到不得已时,更常用的数据库拆分手段是业务分库,将不同的业务数据库部署在不同的物理服务器上。

分布式组件
- 场景:随着网站业务越来越复杂,对数据存储和检索的需求也越来越复杂。
- 特征:系统引入 NoSQL 数据库及搜索引擎。
- 描述:NoSQL 数据库及搜索引擎对可伸缩的分布式特性具有更好的支持。应用服务器通过统一数据访问模块访问各种数据,减轻应用程序管理诸多数据源的麻烦。

微服务
- 场景:随着业务越拆越小,存储系统越来越庞大,应用系统整体复杂程度呈指数级上升,部署维护越来越困难。由于所有应用要和所有数据库系统连接,最终导致数据库连接资源不足,拒绝服务。
- 特征:公共业务提取出来,独立部署。由这些可复用的业务连接数据库,通过分布式服务提供共用业务服务。
- 描述:大型网站的架构演化到这里,基本上大多数的技术问题都得以解决,诸如跨数据中心的实时数据同步和具体网站业务相关的问题也都可以组合改进现有技术架构来解决。

分布式指标
分布式系统的目标是提升系统的整体性能和吞吐量,另外还要尽量保证分布式系统的容错性。
由分布式系统的目标很容易得出分布式系统的关键指标:性能、可用性、可扩展性。这些指标,正对应着耳熟能详的分布式系统“三高”特性——高并发、高性能、高可用。
性能(Performance)
性能用于衡量一个系统处理各种任务的能力。
常见的性能指标有:
- 吞吐量(Throughput) - 系统在一定时间内可以处理的任务数。常见的吞吐量指标有:
- QPS - Queries Per Second 的缩写,即每秒查询数。
- TPS - Transactions Per Second 的缩写,即每秒事务数。
- 响应时间(Response Time) - 执行一个请求从开始到最后收到响应数据所花费的总体时间,即从客户端发起请求到收到服务器响应结果的时间。
- 并发数(Concurrency) - 并发数是指系统能同时处理请求的数量,这个也反映了系统的负载能力。并发意味着可以同时进行多个处理。并发在现代编程中无处不在,网络中有多台计算机同时存在,一台计算机上同时运行着多个应用程序。
以上三个指标的关系大致为:
QPS(TPS)= 并发数 / 平均响应时间
并发数 = QPS(TPS) * 平均响应时间可用性(Availability)
可用性:指的是系统在面对各种异常时可以正确提供服务的能力。
系统的可用性可以用系统停止服务的时间与总的时间之比衡量。
行业内一般用几个 9 表示可用性指标,对应用的可用性程度一般衡量标准有三个 9 到五个 9;一般我们的系统至少要到 4 个 9(99.99%)的可用性才能谈得上高可用。
| 可用性 | 年故障时间 |
|---|---|
| 99.9999% | 32 秒 |
| 99.999% | 5 分 15 秒 |
| 99.99% | 52 分 34 秒 |
| 99.9% | 8 小时 46 分 |
| 99% | 3 天 15 小时 36 分 |
而所谓的高可用,就是:在任何情况下,让服务尽最大可能对外提供服务。
可扩展性(Scalability)
**可扩展性(Scalability)**指的是分布式系统通过扩展集群机器规模提高系统性能 (吞吐、响应时间、 完成时间)、存储容量、计算能力的特性,是分布式系统的特有性质。
系统扩展可以分为垂直扩展、水平扩展。
- 垂直扩展,即提升单机的硬件处理能力,比如 CPU 处理能力,内存容量,磁盘等方面。但是,单机是有性能瓶颈的,一旦触及瓶颈,再想提升,付出的成本和代价会极高。通俗来说,就三个字:得加钱!
- 水平扩展:采用分而治之的思想,通过集群来分担吞吐量。集群中的应用机器(节点)通常被设计成无状态,用户可以请求任何一个节点,这些节点共同分担访问压力。水平扩展有两个要点:
- 集群化、分区化:将一个完整的应用化整为零,如果是无状态应用,可以直接集群化部署;如果是有状态应用,可以将状态数据分区(分片),然后部署到多台机器上。
- 负载均衡:集群化、分区化后,要解决的问题是,请求应该被分发(寻址)到哪台机器上。这就需要通过某种策略来控制分发,这种技术就是负载均衡。
分布式系统分类
分布式技术错综复杂、知识庞杂,且各种技术相互耦合,所以不容易划分层次。
从应用的维度来看,大致可以将分布式系统分为以下四类:
- 分布式计算:解决应用的分布式计算问题。基于分布式计算模式,包括批处理计算、离线计算、在线计算、融合计算等,根据应用类型构建高效智能的分布式计算框架。
- 分布式存储:解决数据的分布式和多元化问题。包括分布式数据库、分布式文件系统、分布式缓存等,支持不同类型的数据的存储和管理。
- 分布式通信:解决进程间的分布式通信问题。通过消息队列、远程调用等方式,实现简单高效的通信。
- 分布式资源管理:解决资源的分布式和异构性问题。将 CPU、内存、IO 等物理资源虚拟化,新城逻辑资源池,以便统一管理。
此外,分布式系统都需要面对一些共性问题,可以视为分布式系统技术的基石:
- 分布式协同 - 解决分布式状态及数据一致性的问题。代表技术:分布式互斥、分布式共识、分布式选举、分布式选举等。
- 分布式调度 - 解决分布式系统资源、请求分配调度的问题。代表技术:服务注册和发现、服务路由、负载均衡、流量控制等。
- 分布式容错 - 解决分布式系统中故障分析、处理的问题,保证系统整体可靠性。代表技术:链路追踪、故障隔离、故障转移等。
- 分布式部署 - 解决分布式系统部署问题。代表技术:CI/CD、容器化等。
分布式系统的挑战
当程序运行在单机上时,通常会以一种可预测的方式运行:要么正常,要么异常。
一旦程序运行在多台机器上时,面临的场景就会变得复杂而难以预料。在分布式系统中,系统的某些部分可能会出现不可预知的故障,这被称为部分失效(partial failure)。问题的难点就在于部分失效是不确定性的。你甚至不确定请求是否成功了,因为消息通过网络传播的时间也是不确定的!这种不确定性和部分失效的可能性,使得分布式系统难以工作。
扩展阅读:The Eight Fallacies of Distributed Computing - Tech Talk 一文中提出了分布式系统新手常有的 8 种误区。
为什么我们要深刻地认识这 8 个错误?这是因为,我们需要清楚地认识到——在分布式系统中,故障是不可避免的。因此,如果要构建一个可靠的分布式系统,就必须要建立容错机制。很可能大部分组件在大部分时间都正常工作。然而,迟早会有一部分系统出现故障,软件必须以某种方式处理。故障处理必须是软件设计的一部分,并且作为软件的运维,你需要知道在发生故障的情况下,软件可能会表现出怎样的行为。
不可靠的网络
互联网以及大多数数据中心的内部网络(通常是以太网)都是异步网络。当通过网络发送数据包时,数据包可能会丢失或者延迟;同样,回复也可能会丢失或延迟。所以如果没有收到回复,并不能确定消息是否发送成功。传输的过程中,可能有各种各样的问题:
- 请求可能已经丢失(可能是被拔掉了网线)。
- 请求可能正在某个队列中等待,无法马上发送(可能是网络或接收方已经超负荷)。
- 远程接收节点可能已经失效(可能是崩愤或关机)。
- 远程节点可能暂时无法响应(例如正在运行长时间的垃圾回收)。
- 远程接收节点已经完成了请求处理,但回复却在网络中丢失(例如网络交换机配置错误)。
- 远程接收节点已经完成了请求处理,但回复却被延迟处理(例如网络或者发送者的机器过载)。

在大多数情况下,系统并没有准确判断节点是否发生故障的机制。因此,分布式系统中,一般通过超时检测来判断远程节点是否可用。但是,超时无法区分网络和节点故障,且可变的网络延迟有时会导致节点被误认为发生崩溃。
超时检测的一个关键点是超时大小的设置:
超时时间如果设置过大,意味着等待时间更久,才能判定节点失效(在此期间,用户只能等待或拿到错误信息)。
超时时间如果设置过小,虽然可以更快检测故障,但增加了误判的可能——节点可能实际上是活着的。当一个节点被宣告为失效,其承担的职责要交给到其他节点,这个过程会给其他节点以及网络带来额外负担,特别是如果此时系统已经处于高负荷状态。
对此,可以先设置一个经验值,然后通过实验逐步调整:先在多台机器上,多次测量往返时间,以确定延迟的大概范围;然后结合应用特点,在故障检测与过早超时风险之间选择一个合适的中间值。更好的做法是:持续测量响应时间及其变化(抖动),然后根据最新的响应时间分布来动态调整。
不可靠的时钟
时钟和计时非常重要。有许多应用程序以各种方式依赖于时钟,例如:
- 某个请求是否超时了?
- 某项服务的 99 %的响应时间是多少?
- 在过去的五分钟内,服务平均每秒处理多少个查询?
- 用户在我们的网站上浏览花了多段时间?
- 这篇文章什么时候发表?
- 在什么时间发送提醒邮件?
- 这个缓存条目何时过期?
- 日志文件中错误消息的时间戳是多少?
在分布式系统中,时间总是件棘手的问题,由于跨节点通信不可能即时完成,消息经由网络从一台机器到另一台机器总是需要花费时间。收到消息的时间应该晚于发送的时间,但是由于网络的不确定延迟,精确测量面临着很多挑战。这些情况使得多节点通信时很难确定事情发生的先后顺序。
为了保证每台机器的时间同步,最常用的机制是 网络时间协议(Network Time Protocol, NTP),它可以根据一组专门的时间服务器来调整本地时间。需要注意的是,即使使用了 NTP 进行时间同步,但是依然会存在一些误差:一方面受限于 NTP 本身的同步精度,此外还受限于网络通信的延迟。
如果想要保证时序,另一种方案是采用逻辑时钟。**逻辑时钟(logic clock)**是基于递增计数器,对于排序事件来说是更安全的选择。逻辑时钟仅测量事件的相对顺序(无论一个事件发生在另一个事件之前还是之后)。
分布式系统的特性
分布式系统具有以下核心特性:
| 特性 | 说明 |
|---|---|
| 透明性 | 用户感知不到系统的分布式特性,使用时如同使用单机系统。包括访问透明、位置透明、迁移透明、复制透明、并发透明、故障透明等 |
| 可扩展性 | 通过增加节点即可线性提升系统能力,包括垂直扩展(提升单机配置)和水平扩展(增加节点数量) |
| 高可用性 | 系统中部分节点故障不影响整体服务,通过冗余部署、故障转移等机制保障 |
| 容错性 | 系统能在部分节点故障、网络分区等异常情况下继续运行 |
| 一致性 | 多个数据副本之间保持一致的状态,根据强度分为强一致性、弱一致性、最终一致性 |
| 并发性 | 系统中多个节点可同时处理请求,通过并发控制机制协调 |
应用场景
分布式系统在以下场景中被广泛应用:
- 互联网应用:电商、社交、视频等高并发、大数据量的应用,需要分布式架构支撑海量用户访问
- 大数据处理:Hadoop、Spark、Flink 等分布式计算框架,用于处理 TB/PB 级数据
- 云原生应用:基于 Kubernetes、Service Mesh 等技术构建的微服务架构
- 金融系统:银行、证券、支付等对数据一致性、可用性要求极高的系统
- 物联网(IoT):海量设备接入、数据采集和处理场景
- 分布式存储:HDFS、Ceph、MinIO 等用于存储海量数据的系统
最佳实践
案例 1:电商系统的高可用架构设计
电商系统在促销活动期间面临突发流量冲击,需要通过分布式架构保障系统稳定性。以下是一个典型的高可用架构实践:
// 使用 Sentinel 进行流量控制,保护核心服务
import com.alibaba.csp.sentinel.annotation.SentinelResource;
import com.alibaba.csp.sentinel.slots.block.BlockException;
import org.springframework.stereotype.Service;
@Service
public class OrderService {
/**
* 下单接口,配置流控规则和降级策略
*/
@SentinelResource(value = "createOrder",
blockHandler = "createOrderBlockHandler",
fallback = "createOrderFallback")
public OrderResult createOrder(OrderRequest request) {
// 1. 参数校验
validateOrderRequest(request);
// 2. 扣减库存(分布式事务)
inventoryService.deduct(request.getProductId(), request.getQuantity());
// 3. 创建订单
Order order = orderRepository.save(buildOrder(request));
// 4. 异步发送消息(解耦)
messageProducer.send("order-topic", new OrderCreatedEvent(order.getId()));
return OrderResult.success(order);
}
/**
* 限流处理:当请求超过阈值时调用
*/
public OrderResult createOrderBlockHandler(OrderRequest request, BlockException ex) {
return OrderResult.fail("系统繁忙,请稍后重试");
}
/**
* 降级处理:当服务出现异常时调用
*/
public OrderResult createOrderFallback(OrderRequest request, Throwable ex) {
log.error("创建订单失败,执行降级", ex);
// 记录失败请求,后续补偿
failureRecorder.record(request);
return OrderResult.fail("系统暂时不可用,已记录您的请求");
}
}配合 YAML 配置:
# application.yml - Sentinel 流控配置
spring:
cloud:
sentinel:
transport:
dashboard: localhost:8080
eager: true
datasource:
flow:
nacos:
server-addr: localhost:8848
dataId: ${spring.application.name}-flow-rules
groupId: SENTINEL_GROUP
rule-type: flow
degrade:
nacos:
server-addr: localhost:8848
dataId: ${spring.application.name}-degrade-rules
groupId: SENTINEL_GROUP
rule-type: degrade案例 2:基于 Redis Cluster 的分布式缓存
使用 Redis Cluster 构建分布式缓存,解决单机 Redis 的容量和性能瓶颈:
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.concurrent.TimeUnit;
@Component
public class DistributedCache {
@Resource
private StringRedisTemplate redisTemplate;
/**
* 写入缓存(带过期时间)
*/
public void set(String key, String value, long timeout, TimeUnit unit) {
redisTemplate.opsForValue().set(key, value, timeout, unit);
}
/**
* 读取缓存,缓存未命中时回源数据库
*/
public String getWithFallback(String key, java.util.function.Supplier<String> dbLoader) {
String value = redisTemplate.opsForValue().get(key);
if (value != null) {
return value;
}
// 缓存未命中,回源数据库
value = dbLoader.get();
if (value != null) {
// 防止缓存穿透,空值也缓存(较短的过期时间)
redisTemplate.opsForValue().set(key, value, 30, TimeUnit.MINUTES);
}
return value;
}
/**
* 使用 Redisson 分布式锁防止缓存击穿
*/
public String getWithLock(String key, java.util.function.Supplier<String> dbLoader) {
String value = redisTemplate.opsForValue().get(key);
if (value != null) {
return value;
}
// 获取分布式锁,只允许一个线程回源
String lockKey = "lock:" + key;
try {
Boolean locked = redisTemplate.opsForValue()
.setIfAbsent(lockKey, "1", 10, TimeUnit.SECONDS);
if (Boolean.TRUE.equals(locked)) {
// 双重检查
value = redisTemplate.opsForValue().get(key);
if (value != null) {
return value;
}
value = dbLoader.get();
if (value != null) {
redisTemplate.opsForValue().set(key, value, 1, TimeUnit.HOURS);
}
return value;
} else {
// 等待后重试
Thread.sleep(50);
return getWithLock(key, dbLoader);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("获取缓存被中断", e);
} finally {
redisTemplate.delete(lockKey);
}
}
}Redis Cluster 配置(redis.conf):
# 启用集群模式
cluster-enabled yes
# 集群配置文件
cluster-config-file nodes-6379.conf
# 节点超时时间(毫秒)
cluster-node-timeout 15000
# 开启 AOF 持久化
appendonly yes
appendfsync everysec案例 3:基于 Kafka 的分布式消息处理
使用 Kafka 实现服务间的异步解耦,削峰填谷:
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
@Component
public class OrderMessageConsumer {
/**
* 消费订单消息,实现幂等性处理
*/
@KafkaListener(topics = "order-topic", groupId = "order-consumer-group")
public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) {
try {
OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);
// 幂等性检查:防止重复消费
if (idempotentRepository.exists(event.getEventId())) {
log.warn("消息已处理,跳过。eventId={}", event.getEventId());
ack.acknowledge();
return;
}
// 处理业务逻辑
processOrderEvent(event);
// 标记已处理
idempotentRepository.markProcessed(event.getEventId());
// 手动提交 offset
ack.acknowledge();
} catch (Exception e) {
log.error("处理订单消息失败,record={}", record, e);
// 进入死信队列,后续人工处理
sendToDeadLetterQueue(record, e);
ack.acknowledge();
}
}
}Kafka 生产者和消费者配置:
# Kafka 生产者配置
spring:
kafka:
producer:
bootstrap-servers: kafka1:9092,kafka2:9092,kafka3:9092
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
acks: all # 所有副本确认后才算发送成功
retries: 3 # 发送失败重试次数
properties:
max.in.flight.requests.per.connection: 5
enable.idempotence: true # 开启幂等性
consumer:
bootstrap-servers: kafka1:9092,kafka2:9092,kafka3:9092
group-id: order-consumer-group
auto-offset-reset: earliest
enable-auto-commit: false # 关闭自动提交,手动管理 offset
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
listener:
ack-mode: manual_immediate # 手动立即提交常见问题
问题 1:分布式系统中调用超时如何排查
问题描述:微服务架构下,某个接口偶尔出现调用超时,但下游服务本身没有明显的性能问题。
原因分析:
- 网络抖动导致延迟增加
- 下游服务 GC(如 Full GC)导致暂停
- 线程池满,请求排队等待
- 超时时间设置不合理
- 跨可用区调用导致延迟
解决方案:
import io.github.resilience4j.retry.Retry;
import io.github.resilience4j.retry.RetryConfig;
import java.time.Duration;
// 使用 Resilience4j 配置重试和超时
public class RemoteCallService {
private final Retry retry;
public RemoteCallService() {
RetryConfig config = RetryConfig.custom()
.maxAttempts(3) // 最大重试次数
.waitDuration(Duration.ofMillis(500)) // 重试间隔
.retryOnException(e -> e instanceof TimeoutException)
.build();
this.retry = Retry.of("remoteCall", config);
}
// 设置合理的超时时间,并加入监控
@Timed(value = "remote.call.duration", description = "远程调用耗时")
public String callRemoteService(String param) {
return Retry.decorateSupplier(retry, () -> {
// 使用 OkHttp 设置连接和读取超时
OkHttpClient client = new OkHttpClient.Builder()
.connectTimeout(2, TimeUnit.SECONDS) // 连接超时
.readTimeout(5, TimeUnit.SECONDS) // 读取超时
.writeTimeout(5, TimeUnit.SECONDS) // 写入超时
.build();
Request request = new Request.Builder()
.url("http://downstream-service/api/process?param=" + param)
.build();
try (Response response = client.newCall(request).execute()) {
if (!response.isSuccessful()) {
throw new RuntimeException("调用失败: " + response.code());
}
return response.body().string();
} catch (IOException e) {
throw new RuntimeException("远程调用异常", e);
}
}).get();
}
}问题 2:分布式事务数据不一致
问题描述:在跨服务的数据操作中(如订单服务创建订单 + 库存服务扣减库存),由于网络故障或服务异常,导致数据不一致:订单创建成功但库存未扣减。
原因分析:
- 使用了非原子的跨服务调用,无法保证两个操作同时成功或失败
- 网络分区导致消息丢失
- 服务异常时没有完善的补偿机制
解决方案:使用 Seata 的 AT 模式实现分布式事务:
import io.seata.spring.annotation.GlobalTransactional;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
@Service
public class OrderTransactionService {
@Resource
private OrderService orderService;
@Resource
private InventoryFeignClient inventoryFeignClient;
/**
* 使用 Seata 全局事务,保证订单创建和库存扣减的原子性
*/
@GlobalTransactional(name = "createOrderTx", timeoutMills = 60000, rollbackFor = Exception.class)
public void createOrderWithTx(OrderRequest request) {
// 1. 创建订单(本地事务)
Order order = orderService.create(request);
// 2. 扣减库存(远程调用,由 Seata 协调)
boolean success = inventoryFeignClient.deduct(request.getProductId(), request.getQuantity());
if (!success) {
throw new RuntimeException("库存扣减失败");
}
// 3. 如果后续步骤失败,Seata 会自动回滚前面的操作
}
}Seata 配置(application.yml):
seata:
enabled: true
application-id: order-service
tx-service-group: my_tx_group
service:
vgroup-mapping:
my_tx_group: default
grouplist:
default: 127.0.0.1:8091
registry:
type: nacos
nacos:
server-addr: 127.0.0.1:8848
namespace: ""
group: SEATA_GROUP
config:
type: nacos
nacos:
server-addr: 127.0.0.1:8848
namespace: ""
group: SEATA_GROUP问题 3:分布式锁失效导致并发问题
问题描述:使用基于 Redis 的分布式锁时,由于锁过期但业务未执行完,导致多个节点同时获取到锁,出现并发问题。
原因分析:
- 锁的过期时间设置过短,业务还没执行完锁就过期了
- 节点 A 获取锁后因 GC 暂停,锁过期后节点 B 获取锁,A 恢复后继续执行,造成并发
- 没有正确实现锁的可重入性
解决方案:使用 Redisson 的分布式锁,支持自动续期(Watchdog 机制):
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.concurrent.TimeUnit;
@Component
public class DistributedLockService {
@Resource
private RedissonClient redissonClient;
/**
* 使用 Redisson 分布式锁,支持自动续期
* @param lockKey 锁的 key
* @param task 需要加锁执行的任务
*/
public <T> T executeWithLock(String lockKey, java.util.function.Supplier<T> task) {
RLock lock = redissonClient.getLock(lockKey);
try {
// 尝试加锁,waitTime=10秒,leaseTime=-1 表示启用 Watchdog 自动续期
boolean acquired = lock.tryLock(10, -1, TimeUnit.SECONDS);
if (!acquired) {
throw new RuntimeException("获取分布式锁失败: " + lockKey);
}
return task.get();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("获取锁被中断", e);
} finally {
// 确保只有持有锁的线程才能释放
if (lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}
}Redisson 配置:
import org.redisson.Redisson;
import org.redisson.config.Config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RedissonConfig {
@Bean(destroyMethod = "shutdown")
public RedissonClient redissonClient() {
Config config = new Config();
config.useClusterServers()
.addNodeAddress(
"redis://127.0.0.1:7000",
"redis://127.0.0.1:7001",
"redis://127.0.0.1:7002"
)
.setConnectTimeout(3000)
.setTimeout(3000)
.setRetryAttempts(3)
.setRetryInterval(500);
// Watchdog 超时时间,默认 30 秒
config.setLockWatchdogTimeout(30000);
return Redisson.create(config);
}
}参考资料
- 《数据密集型应用系统设计》 - 这可能是目前最好的分布式存储书籍,强力推荐【进阶】
- The Eight Fallacies of Distributed Computing - Tech Talk - 分布式系统新手常犯的 8 个错误,并探讨了其会带来的影响。
- Distributed Systems for Fun and Profit - 一本学习小册,涵盖了分布式系统中的关键问题,包括时间的作用和不同的复制策略。
- A Note on Distributed Systems - 这是一篇经典的论文,讲述了为什么在分布式系统中,远程交互不能像本地对象那样进行。
- Time, Clocks, and the Ordering of Events in a Distributed System,译文,解读 - Lamport 介绍 happened before、偏序关系(partial ordering)、逻辑时钟(Logical Clocks)概念,提出解决分布式系统中区分事件发生的时序问题的方法。
- Virtual Time and Global States of Distributed Systems,解读 - 逻辑时钟无法描述事件的因果关系。本文提出了向量时钟,这种算法利用了向量这种数据结构将全局各个进程的逻辑时间戳广播给各个进程,通过向量时间戳就能够比较任意两个事件的因果关系。