分布式高并发数据一致性:从理论到实战,一篇讲透!
引言:从奶茶店的故事说起
假设你开了一家网红奶茶店,生意火爆。今天是周六,顾客排起了长队。这时候,你会遇到几个问题:
-
只有一个收银员:收银员既要收钱,又要记库存,忙不过来。这就像单体应用。
-
你增加了收银员:现在有两个收银员,但他们都用一个账本记录库存。这就像多线程并发。
-
你开了分店:总店和分店共享库存,但沟通靠打电话。这就像分布式系统。
-
节假日高峰期:两家店同时有100人排队,都来买同一款限量奶茶。这就像高并发。
-
最头疼的问题:总店卖了最后一杯奶茶,还没来得及告诉分店,分店也卖了这最后一杯。结果超卖了,顾客生气。这就是数据不一致。
今天,我就用最通俗的语言,带你彻底理解这个复杂的分布式高并发数据一致性问题。无论你是刚入门的小白,还是想深入理解的老手,这篇文章都会给你清晰的答案。
第一章:基础篇 - 这到底是什么问题?
1.1 三个核心概念的通俗解释
什么是分布式系统?
// 原来:一家奶茶店搞定所有事情(单体系统)
public class SingleTeaShop {
public void sellMilkTea() {
collectMoney();
recordSale();
updateInventory(); // 所有操作在一台电脑上完成
makeTea();
}
}
// 现在:开了三家分店(分布式系统)
public class DistributedTeaShops {
@Service("mainShop") // 在北京
class MainShop {
void collectAndRecord() { /* 处理订单 */ }
}
@Service("shanghaiShop") // 在上海,不同服务器
class ShanghaiShop {
void updateInventory() { /* 管理库存 */ }
}
@Service("guangzhouShop") // 在广州,不同服务器
class GuangzhouShop {
void makeTea() { /* 生产制作 */ }
}
// 卖一杯奶茶需要三家店协作
public void sellMilkTea() {
mainShop.collectAndRecord(); // 网络调用
shanghaiShop.updateInventory(); // 网络调用
guangzhouShop.makeTea(); // 网络调用
}
}
什么是高并发?
想象一下:
-
平常:每分钟有5个人买奶茶
-
节假日:每分钟有500人同时抢购(这就是高并发)
-
双十一:每分钟有50000人同时下单(这就是超高并发)
什么是数据一致性?
继续奶茶店的例子:
-
一致的情况:总店显示奶茶还剩100杯,分店也显示100杯
-
不一致的情况:总店显示还剩1杯,分店却显示还有10杯(实际上只有1杯)
-
超卖的后果:两家店都卖了这最后一杯,但顾客只能拿到一杯,另一杯要退款道歉
1.2 数据为什么会"打架"?
情况1:两个人同时修改数据(并发问题)
public class ConcurrencyProblem {
static int teaInventory = 10; // 还剩10杯
public static void main(String[] args) throws InterruptedException {
// 顾客A要买5杯
Thread customerA = new Thread(() -> {
if (teaInventory >= 5) {
try {
Thread.sleep(100); // 模拟处理时间
} catch (InterruptedException e) {
e.printStackTrace();
}
teaInventory = teaInventory - 5; // A以为库存够,开始扣减
System.out.println("顾客A购买5杯,库存: " + teaInventory);
}
});
// 顾客B要买6杯
Thread customerB = new Thread(() -> {
if (teaInventory >= 6) {
try {
Thread.sleep(100); // 模拟处理时间
} catch (InterruptedException e) {
e.printStackTrace();
}
teaInventory = teaInventory - 6; // B也以为库存够,开始扣减
System.out.println("顾客B购买6杯,库存: " + teaInventory);
}
});
// 两个人同时执行
customerA.start();
customerB.start();
customerA.join();
customerB.join();
// 结果:库存可能变成 -1!这就是超卖
System.out.println("最终库存: " + teaInventory);
}
}
情况2:网络延迟(像微信消息延迟)
-
总店卖出了最后一杯,库存变为0
-
打电话告诉分店:"没货了!"
-
但电话信号不好,分店10分钟后才收到消息
-
这10分钟内,分店又卖出了几杯"不存在"的奶茶
情况3:系统故障(像写作业时停电)
public class SystemFailure {
public void transferMoney(String fromAccount, String toAccount, int amount) {
// 第一步:从张三账户扣钱
deductFromAccount(fromAccount, amount);
// 这时候银行系统崩溃了!停电了!
// 程序停止运行...
// 第二步:给李四账户加钱(永远执行不到)
addToAccount(toAccount, amount);
}
private void deductFromAccount(String account, int amount) {
// 扣款逻辑
}
private void addToAccount(String account, int amount) {
// 加款逻辑
}
}
// 结果:张三的钱没了,李四没收到钱,钱凭空消失了!
第二章:解决方案进化史
2.1 原始方案:单店模式(单体应用)
// 最简单的方法:只有一个收银员,一本账
public class SingleShopSolution {
private int inventory = 100; // 库存100杯
private final Object lock = new Object(); // 锁对象
public boolean sellTea(int quantity) {
synchronized (lock) { // 加锁:像厕所门栓,一次只进一个人
if (inventory >= quantity) {
// 模拟一些操作时间
try { Thread.sleep(100); } catch (Exception e) {}
inventory -= quantity;
System.out.println("卖出" + quantity + "杯,剩余:" + inventory);
return true;
}
return false;
}
}
// 测试并发
public static void main(String[] args) throws InterruptedException {
SingleShopSolution shop = new SingleShopSolution();
ExecutorService executor = Executors.newFixedThreadPool(10);
for (int i = 0; i < 100; i++) {
executor.submit(() -> {
shop.sellTea(1);
});
}
executor.shutdown();
executor.awaitTermination(10, TimeUnit.SECONDS);
System.out.println("最终库存: " + shop.inventory);
}
}
2.2 基础分布式方案:分店模式(分布式锁)
方案1:电话确认法(分布式锁)
public class PhoneConfirmSolution {
private Jedis redis = new Jedis("localhost", 6379);
public boolean sellTea(String shopName, int quantity) {
String lockKey = "tea_lock";
String requestId = UUID.randomUUID().toString();
try {
// 1. 尝试获取锁(设置10秒超时)
boolean locked = redis.setnx(lockKey, requestId) == 1;
if (locked) {
redis.expire(lockKey, 10); // 设置10秒过期
} else {
return false; // 锁被占用,稍后重试
}
// 2. 查询库存
int currentStock = Integer.parseInt(redis.get("tea_inventory") == null ? "100" : redis.get("tea_inventory"));
if (currentStock >= quantity) {
// 3. 扣减库存
redis.set("tea_inventory", String.valueOf(currentStock - quantity));
return true;
}
return false;
} finally {
// 4. 必须释放锁
if (requestId.equals(redis.get(lockKey))) {
redis.del(lockKey);
}
}
}
}
方案2:小黑板法(乐观锁)
public class OptimisticLockSolution {
private Jedis redis = new Jedis("localhost", 6379);
public boolean sellTea(int quantity) {
int retryCount = 0;
while (retryCount < 3) { // 最多重试3次
// 1. 读取当前库存和版本
int currentStock = Integer.parseInt(redis.get("tea_inventory") == null ? "100" : redis.get("tea_inventory"));
int currentVersion = Integer.parseInt(redis.get("tea_version") == null ? "1" : redis.get("tea_version"));
if (currentStock < quantity) {
return false; // 库存不足
}
// 2. 使用WATCH/MULTI/EXEC实现乐观锁
redis.watch("tea_inventory", "tea_version");
Transaction transaction = redis.multi();
transaction.set("tea_inventory", String.valueOf(currentStock - quantity));
transaction.incr("tea_version");
List<Object> results = transaction.exec();
if (results != null && !results.isEmpty()) {
return true; // 更新成功
}
// 3. 更新失败(别人先修改了),稍后重试
retryCount++;
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
return false;
}
}
2.3 两阶段提交(2PC) - 分布式事务的基础协议
// 2PC协调者
public class TwoPhaseCommitCoordinator {
private List<Participant> participants = new ArrayList<>();
public boolean executeTransaction() {
// 第一阶段:准备阶段(投票阶段)
List<Boolean> prepareResults = new ArrayList<>();
for (Participant participant : participants) {
try {
boolean prepared = participant.prepare();
prepareResults.add(prepared);
} catch (Exception e) {
prepareResults.add(false);
}
}
// 检查所有参与者是否都准备成功
boolean allPrepared = prepareResults.stream().allMatch(result -> result);
// 第二阶段:提交/回滚阶段
if (allPrepared) {
// 所有参与者都准备成功,提交事务
for (Participant participant : participants) {
participant.commit();
}
return true;
} else {
// 有参与者准备失败,回滚事务
for (Participant participant : participants) {
participant.rollback();
}
return false;
}
}
// 参与者接口
public interface Participant {
boolean prepare(); // 准备阶段:预留资源
void commit(); // 提交阶段:正式提交
void rollback(); // 回滚阶段:释放资源
}
// 数据库参与者示例
public class DatabaseParticipant implements Participant {
private Connection connection;
public DatabaseParticipant(Connection connection) {
this.connection = connection;
}
@Override
public boolean prepare() {
try {
// 1. 记录prepare日志
PreparedStatement stmt = connection.prepareStatement(
"INSERT INTO transaction_log (tx_id, status) VALUES (?, 'PREPARED')"
);
stmt.setString(1, UUID.randomUUID().toString());
stmt.executeUpdate();
// 2. 锁定相关资源(这里简化处理)
connection.setAutoCommit(false);
return true;
} catch (SQLException e) {
return false;
}
}
@Override
public void commit() {
try {
// 执行真正的提交
connection.commit();
// 更新事务状态
PreparedStatement stmt = connection.prepareStatement(
"UPDATE transaction_log SET status = 'COMMITTED' WHERE tx_id = ?"
);
stmt.setString(1, getCurrentTxId());
stmt.executeUpdate();
connection.setAutoCommit(true);
} catch (SQLException e) {
// 记录错误
}
}
@Override
public void rollback() {
try {
// 执行回滚
connection.rollback();
// 更新事务状态
PreparedStatement stmt = connection.prepareStatement(
"UPDATE transaction_log SET status = 'ROLLBACKED' WHERE tx_id = ?"
);
stmt.setString(1, getCurrentTxId());
stmt.executeUpdate();
connection.setAutoCommit(true);
} catch (SQLException e) {
// 记录错误
}
}
private String getCurrentTxId() {
// 获取当前事务ID
return "";
}
}
// 业务使用示例
public class OrderService {
private TwoPhaseCommitCoordinator coordinator;
public void createOrder(Order order) {
// 添加参与者
coordinator.addParticipant(new DatabaseParticipant(orderDbConnection));
coordinator.addParticipant(new DatabaseParticipant(inventoryDbConnection));
coordinator.addParticipant(new DatabaseParticipant(paymentDbConnection));
// 执行2PC事务
boolean success = coordinator.executeTransaction();
if (!success) {
throw new RuntimeException("订单创建失败,事务回滚");
}
}
}
}
2PC的现实比喻:
想象一场需要所有人到场的婚礼:
-
阶段一(准备阶段):司仪问每个人"你愿意吗?"
-
阶段二(提交阶段):
-
如果所有人都说"愿意" → 司仪宣布"礼成"
-
如果有人说不愿意 → 司仪宣布"婚礼取消"
-
2PC的优缺点:
-
优点:强一致性保证
-
缺点:
-
同步阻塞:所有参与者都要等待
-
单点故障:协调者挂了,所有参与者卡住
-
数据不一致:协调者发送commit后崩溃,部分参与者没收到
-
2.4 高级分布式方案:现代互联网方案
方案3:预定登记法(TCC模式)
TCC = Try-Confirm-Cancel(尝试-确认-取消)
public class TCCExample {
// 第一阶段:Try(尝试)
public boolean tryReserve(String orderId, String productId, int quantity) {
// 1. 检查库存(不实际扣减)
if (!checkInventory(productId, quantity)) {
return false;
}
// 2. 预留资源(冻结库存)
boolean reserved = freezeInventory(productId, quantity, orderId);
// 3. 记录TCC事务日志
if (reserved) {
saveTccLog(orderId, "TRY", "INVENTORY_RESERVED");
}
return reserved;
}
// 第二阶段:Confirm(确认)
public boolean confirmReservation(String orderId) {
// 1. 查询事务日志
TccLog log = getTccLog(orderId);
if (log == null || !"TRY".equals(log.getStatus())) {
return false;
}
// 2. 真正扣减库存
boolean confirmed = deductInventory(log.getProductId(), log.getQuantity());
// 3. 更新事务状态
if (confirmed) {
updateTccLog(orderId, "CONFIRMED");
}
return confirmed;
}
// 第三阶段:Cancel(取消)
public boolean cancelReservation(String orderId) {
// 1. 查询事务日志
TccLog log = getTccLog(orderId);
if (log == null || !"TRY".equals(log.getStatus())) {
return false;
}
// 2. 解冻库存
boolean cancelled = unfreezeInventory(log.getProductId(), log.getQuantity(), orderId);
// 3. 更新事务状态
if (cancelled) {
updateTccLog(orderId, "CANCELLED");
}
return cancelled;
}
// TCC事务执行器
public void executeTccTransaction(Order order) {
String orderId = order.getId();
try {
// 第一阶段:Try
if (tryReserve(orderId, order.getProductId(), order.getQuantity())) {
// 第二阶段:Confirm
confirmReservation(orderId);
}
} catch (Exception e) {
// 出现异常,执行Cancel
cancelReservation(orderId);
throw e;
}
}
// 辅助方法
private boolean checkInventory(String productId, int quantity) {
// 检查库存是否足够
return true;
}
private boolean freezeInventory(String productId, int quantity, String orderId) {
// 冻结库存
return true;
}
private boolean deductInventory(String productId, int quantity) {
// 扣减库存
return true;
}
private boolean unfreezeInventory(String productId, int quantity, String orderId) {
// 解冻库存
return true;
}
private void saveTccLog(String orderId, String status, String data) {
// 保存TCC日志
}
private TccLog getTccLog(String orderId) {
// 获取TCC日志
return null;
}
private void updateTccLog(String orderId, String status) {
// 更新TCC日志
}
// TCC日志类
static class TccLog {
private String orderId;
private String status;
private String productId;
private int quantity;
// getters and setters
public String getStatus() { return status; }
public String getProductId() { return productId; }
public int getQuantity() { return quantity; }
}
}
方案4:旅行计划法(SAGA模式)
public class SagaExample {
// 正向操作序列
private List<SagaStep> steps = Arrays.asList(
new SagaStep("bookFlight", this::bookFlight, this::cancelFlight),
new SagaStep("bookHotel", this::bookHotel, this::cancelHotel),
new SagaStep("rentCar", this::rentCar, this::cancelCarRental),
new SagaStep("makePayment", this::makePayment, this::refundPayment)
);
// 执行SAGA事务
public void bookTravel(TravelRequest request) {
List<ExecutedStep> executedSteps = new ArrayList<>();
String sagaId = UUID.randomUUID().toString();
try {
// 按顺序执行每个步骤
for (SagaStep step : steps) {
step.getAction().execute(request, sagaId);
executedSteps.add(new ExecutedStep(step, true, sagaId));
}
// 所有步骤成功,SAGA完成
System.out.println("Travel booked successfully!");
} catch (Exception e) {
System.out.println("Error occurred, starting compensation...");
// 按相反顺序执行补偿操作
Collections.reverse(executedSteps);
for (ExecutedStep executedStep : executedSteps) {
if (executedStep.isSuccess()) {
try {
executedStep.getStep().getCompensation().execute(request, sagaId);
System.out.println("Compensated: " + executedStep.getStep().getName());
} catch (Exception ex) {
System.err.println("Compensation failed for: " + executedStep.getStep().getName());
// 记录日志,需要人工干预
}
}
}
throw e;
}
}
// 正向操作:预订航班
public void bookFlight(TravelRequest request, String sagaId) {
// 预订航班逻辑
System.out.println("Booking flight for: " + request.getUserId());
// 如果失败,抛出异常
}
// 补偿操作:取消航班
public void cancelFlight(TravelRequest request, String sagaId) {
// 取消航班逻辑
System.out.println("Canceling flight for: " + request.getUserId());
}
// 其他正向和补偿操作...
// 内部类定义
class SagaStep {
private String name;
private BiConsumer<TravelRequest, String> action;
private BiConsumer<TravelRequest, String> compensation;
public SagaStep(String name, BiConsumer<TravelRequest, String> action,
BiConsumer<TravelRequest, String> compensation) {
this.name = name;
this.action = action;
this.compensation = compensation;
}
public String getName() { return name; }
public BiConsumer<TravelRequest, String> getAction() { return action; }
public BiConsumer<TravelRequest, String> getCompensation() { return compensation; }
}
class ExecutedStep {
private SagaStep step;
private boolean success;
private String sagaId;
public ExecutedStep(SagaStep step, boolean success, String sagaId) {
this.step = step;
this.success = success;
this.sagaId = sagaId;
}
public SagaStep getStep() { return step; }
public boolean isSuccess() { return success; }
}
class TravelRequest {
private String userId;
private String destination;
private Date travelDate;
// getters and setters
public String getUserId() { return userId; }
}
}
方案5:本地消息表(最实用的最终一致性方案)
public class LocalMessageTableSolution {
@Autowired
private OrderRepository orderRepository;
@Autowired
private MessageRepository messageRepository;
@Autowired
private MQProducer mqProducer;
@Transactional
public Order createOrder(OrderRequest request) {
// 1. 创建订单(主业务)
Order order = new Order();
order.setOrderId(UUID.randomUUID().toString());
order.setUserId(request.getUserId());
order.setProductId(request.getProductId());
order.setQuantity(request.getQuantity());
order.setStatus("CREATED");
orderRepository.save(order);
// 2. 在同一个事务中,保存消息
Message message = new Message();
message.setMessageId(UUID.randomUUID().toString());
message.setTopic("ORDER_CREATED");
message.setContent(JSON.toJSONString(order));
message.setStatus("PENDING");
message.setRetryCount(0);
message.setCreatedTime(new Date());
messageRepository.save(message);
// 3. 事务提交(保证业务和消息的原子性)
return order;
}
// 定时任务:扫描并发送消息
@Scheduled(fixedDelay = 5000)
@Transactional
public void processPendingMessages() {
List<Message> pendingMessages = messageRepository.findByStatus("PENDING");
for (Message message : pendingMessages) {
try {
// 发送到消息队列
boolean sent = mqProducer.send(message.getTopic(), message.getContent());
if (sent) {
// 发送成功,更新状态
message.setStatus("SENT");
message.setSentTime(new Date());
messageRepository.save(message);
} else {
// 发送失败,增加重试次数
handleFailedMessage(message);
}
} catch (Exception e) {
handleFailedMessage(message);
}
}
}
private void handleFailedMessage(Message message) {
message.setRetryCount(message.getRetryCount() + 1);
if (message.getRetryCount() > 3) {
message.setStatus("FAILED");
// 发送告警,需要人工处理
alertManualProcess(message);
}
messageRepository.save(message);
}
// 消费者端(保证幂等性)
@MQListener(topic = "ORDER_CREATED")
public void handleOrderCreated(String messageContent) {
// 1. 解析消息
OrderMessage message = JSON.parseObject(messageContent, OrderMessage.class);
// 2. 检查消息是否已处理(幂等性)
if (messageRepository.existsByMessageId(message.getMessageId())) {
return; // 已处理过,直接返回
}
try {
// 3. 执行业务逻辑(扣减库存等)
inventoryService.deductStock(message.getProductId(), message.getQuantity());
// 4. 记录已处理的消息
Message processedMessage = new Message();
processedMessage.setMessageId(message.getMessageId());
processedMessage.setStatus("PROCESSED");
processedMessage.setProcessedTime(new Date());
messageRepository.save(processedMessage);
} catch (Exception e) {
// 消费失败,抛出异常让MQ重试
throw e;
}
}
// 实体类
@Entity
@Table(name = "order_message")
static class Message {
@Id
private String messageId;
private String topic;
private String content;
private String status;
private Integer retryCount;
private Date createdTime;
private Date sentTime;
private Date processedTime;
// getters and setters
public String getMessageId() { return messageId; }
public void setMessageId(String messageId) { this.messageId = messageId; }
public String getTopic() { return topic; }
public void setTopic(String topic) { this.topic = topic; }
public String getContent() { return content; }
public void setContent(String content) { this.content = content; }
public String getStatus() { return status; }
public void setStatus(String status) { this.status = status; }
public Integer getRetryCount() { return retryCount; }
public void setRetryCount(Integer retryCount) { this.retryCount = retryCount; }
public Date getCreatedTime() { return createdTime; }
public void setCreatedTime(Date createdTime) { this.createdTime = createdTime; }
public Date getSentTime() { return sentTime; }
public void setSentTime(Date sentTime) { this.sentTime = sentTime; }
public Date getProcessedTime() { return processedTime; }
public void setProcessedTime(Date processedTime) { this.processedTime = processedTime; }
}
static class OrderMessage {
private String messageId;
private String orderId;
private String productId;
private Integer quantity;
// getters and setters
public String getMessageId() { return messageId; }
public void setMessageId(String messageId) { this.messageId = messageId; }
public String getProductId() { return productId; }
public void setProductId(String productId) { this.productId = productId; }
public Integer getQuantity() { return quantity; }
public void setQuantity(Integer quantity) { this.quantity = quantity; }
}
}
第三章:理论篇 - 分布式系统的"交通规则"
3.1 CAP定理:分布式系统的"不可能三角"
想象一下城市的交通管理:
C(一致性):所有路口的红绿灯同步
-
所有路口同时变红或变绿
-
绝对不会出现一个路口绿灯,下一个路口红灯
A(可用性):路口永远有信号灯
-
即使停电,也有交警指挥
-
绝对不会出现"信号灯坏了,请绕行"
P(分区容错性):部分路段封闭时,其他路段还能通行
-
一条路修路,其他路还能走
残酷的现实:你只能三选二!
-
CP系统(一致性+分区容错):
-
像银行系统:宁可暂时不能转账,也要保证账目绝对正确
-
网络出问题时,停止服务,保证数据一致
-
-
AP系统(可用性+分区容错):
-
像电商网站:宁可显示旧价格,也要保证能访问
-
网络出问题时,继续服务,返回可能过期的数据
-
-
CA系统? 在分布式系统中不存在!
-
因为分布式系统一定会有网络问题(P)
-
就像城市一定有修路的时候
-
public class CAPTheorem {
enum SystemType {
BANK_SYSTEM {
// 选择CP:保证一致性
@Override
Object handleNetworkPartition() {
// 停止服务:"系统维护中,请稍后再试"
throw new ServiceUnavailableException("系统维护中,请稍后再试");
}
},
ECOMMERCE_SYSTEM {
// 选择AP:保证可用性
@Override
Object handleNetworkPartition() {
// 继续服务,可能返回旧数据
return getCachedData(); // 可能是昨天的价格
}
},
SOCIAL_MEDIA {
// 选择AP:点赞数晚点更新没关系
@Override
Object handleNetworkPartition() {
return getLocalData();
}
};
abstract Object handleNetworkPartition();
}
// 模拟网络分区场景
public static void main(String[] args) {
SystemType system = SystemType.BANK_SYSTEM;
try {
Object result = system.handleNetworkPartition();
System.out.println("系统响应: " + result);
} catch (ServiceUnavailableException e) {
System.out.println("系统不可用: " + e.getMessage());
}
}
static class ServiceUnavailableException extends RuntimeException {
public ServiceUnavailableException(String message) {
super(message);
}
}
}
3.2 BASE理论:实用主义的选择
BASE = Basically Available, Soft state, Eventually consistent
基本可用,软状态,最终一致
对比传统数据库的ACID:
-
ACID:像数学考试,必须100%准确
-
BASE:像日常交流,大体正确就行,允许延迟
public class BASEExample {
// 场景:朋友圈点赞
public void likePost(String postId, String userId) {
try {
// 1. 快速响应,更新缓存
incrementLikeInCache(postId);
// 2. 异步持久化到数据库
asyncPersistLike(postId, userId);
} catch (ServiceDownException e) {
// 3. 基本可用:即使服务挂了,也有降级方案
logLikeLocally(postId, userId);
}
}
private void incrementLikeInCache(String postId) {
// 快速更新缓存
Jedis jedis = new Jedis("localhost", 6379);
jedis.incr("post_like:" + postId);
}
private void asyncPersistLike(String postId, String userId) {
// 异步保存到数据库
CompletableFuture.runAsync(() -> {
try {
saveLikeToDatabase(postId, userId);
} catch (Exception e) {
// 重试机制
retrySaveLike(postId, userId);
}
});
}
private void logLikeLocally(String postId, String userId) {
// 记录到本地,等网络恢复了再同步
System.out.println("记录本地点赞: postId=" + postId + ", userId=" + userId);
}
// BASE特点总结:
// 1. 基本可用:即使部分功能失效,核心功能仍可用
// 2. 软状态:允许数据暂时不一致
// 3. 最终一致:保证数据最终会一致
}
第四章:实战篇 - 设计一个秒杀系统
现在,让我们用这些知识设计一个真正的秒杀系统。假设我们要卖1000瓶茅台,有10万人抢购。
4.1 完整秒杀系统架构
@Service
public class SeckillSystem {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private ProductRepository productRepository;
@Autowired
private OrderRepository orderRepository;
// 1. 限流控制
public SeckillResult seckillWithRateLimit(String userId, String productId) {
// 令牌桶限流
if (!rateLimiter.tryAcquire()) {
return SeckillResult.error("系统繁忙,请稍后重试");
}
// 2. Redis预减库存
Long stock = redisTemplate.opsForValue().decrement("seckill:stock:" + productId);
if (stock < 0) {
// 库存不足,恢复
redisTemplate.opsForValue().increment("seckill:stock:" + productId);
return SeckillResult.error("已售罄");
}
// 3. 发送异步消息
String seckillId = generateSeckillId(userId, productId);
SeckillMessage message = new SeckillMessage(seckillId, userId, productId);
rocketMQTemplate.sendMessageInTransaction(
"SECKILL_TOPIC",
MessageBuilder.withPayload(message).build(),
null
);
return SeckillResult.success("抢购请求已受理", seckillId);
}
// 4. 本地事务:创建订单
@Transactional
public LocalTransactionState createSeckillOrder(Message msg, Object arg) {
SeckillMessage message = JSON.parseObject(new String((byte[]) msg.getPayload()), SeckillMessage.class);
try {
// 4.1 验证重复购买
if (orderRepository.existsByUserIdAndProductId(message.getUserId(), message.getProductId())) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
// 4.2 数据库最终检查(悲观锁)
Product product = productRepository.findForUpdate(message.getProductId());
if (product.getStock() <= 0) {
// 恢复Redis库存
redisTemplate.opsForValue().increment("seckill:stock:" + message.getProductId());
return LocalTransactionState.ROLLBACK_MESSAGE;
}
// 4.3 扣减数据库库存
product.setStock(product.getStock() - 1);
productRepository.save(product);
// 4.4 创建订单
Order order = new Order();
order.setOrderId(generateOrderId());
order.setUserId(message.getUserId());
order.setProductId(message.getProductId());
order.setType("SECKILL");
order.setStatus("CREATED");
order.setCreateTime(new Date());
orderRepository.save(order);
// 4.5 发送订单创建成功消息
rocketMQTemplate.send("ORDER_CREATED", order);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
// 5. 定时对账(最终一致性保障)
@Scheduled(cron = "0 */5 * * * ?")
@Transactional
public void reconcileStock() {
List<Product> seckillProducts = productRepository.findSeckillProducts();
for (Product product : seckillProducts) {
Long redisStock = (Long) redisTemplate.opsForValue().get("seckill:stock:" + product.getId());
Long dbStock = product.getStock();
if (!Objects.equals(redisStock, dbStock)) {
// 发送告警
alertService.sendAlert("库存不一致告警",
"productId=" + product.getId() +
", redisStock=" + redisStock +
", dbStock=" + dbStock);
// 自动修复:以数据库为准
redisTemplate.opsForValue().set("seckill:stock:" + product.getId(), dbStock);
}
}
}
// 6. 订单超时取消
@Scheduled(cron = "0 */1 * * * ?")
@Transactional
public void cancelTimeoutOrders() {
List<Order> timeoutOrders = orderRepository.findTimeoutOrders(15); // 15分钟未支付
for (Order order : timeoutOrders) {
try {
// 取消订单
order.setStatus("CANCELLED");
order.setCancelTime(new Date());
orderRepository.save(order);
// 恢复库存
productRepository.incrementStock(order.getProductId(), 1);
// 更新Redis库存
redisTemplate.opsForValue().increment("seckill:stock:" + order.getProductId());
} catch (Exception e) {
log.error("取消超时订单失败: orderId=" + order.getOrderId(), e);
}
}
}
// 辅助方法
private String generateSeckillId(String userId, String productId) {
return userId + "_" + productId + "_" + System.currentTimeMillis();
}
private String generateOrderId() {
return "ORD" + System.currentTimeMillis() + RandomUtils.nextInt(1000, 9999);
}
// 实体类
static class SeckillMessage {
private String seckillId;
private String userId;
private String productId;
public SeckillMessage(String seckillId, String userId, String productId) {
this.seckillId = seckillId;
this.userId = userId;
this.productId = productId;
}
// getters and setters
public String getSeckillId() { return seckillId; }
public String getUserId() { return userId; }
public String getProductId() { return productId; }
}
static class SeckillResult {
private boolean success;
private String message;
private String seckillId;
public static SeckillResult success(String message, String seckillId) {
SeckillResult result = new SeckillResult();
result.success = true;
result.message = message;
result.seckillId = seckillId;
return result;
}
public static SeckillResult error(String message) {
SeckillResult result = new SeckillResult();
result.success = false;
result.message = message;
return result;
}
}
}
4.2 系统架构总结
秒杀系统架构层次:
1. 接入层:Nginx负载均衡 + 网关限流
2. 服务层:多实例部署 + 本地缓存
3. 缓存层:Redis集群(库存预热 + 预扣减)
4. 消息层:RocketMQ集群(异步削峰)
5. 数据层:MySQL分库分表 + 读写分离
6. 监控层:Prometheus + Grafana + ELK
第五章:方案选择指南
5.1 根据业务场景选择
| 业务场景 | 推荐方案 | 原因 | 生活比喻 |
|---|---|---|---|
| 银行转账 | TCC + 对账 | 钱不能错,必须强一致 | 像银行柜台,必须100%准确 |
| 电商下单 | 本地消息表 | 允许短暂不一致,要高性能 | 像网购,下单成功,稍后发货 |
| 秒杀抢购 | Redis预扣+异步 | 超高并发,快速响应 | 像双十一抢购,先抢到再说 |
| 社交点赞 | 最终一致性 | 少几个赞没关系 | 像朋友圈,晚点更新没关系 |
| 旅行预订 | SAGA模式 | 长流程,需要补偿 | 像旅行团,一项失败要取消全部 |
| 配置管理 | 分布式锁 | 配置必须一致 | 像文件共享,只能一个人编辑 |
| 传统系统 | 2PC | 强一致性要求 | 像婚礼仪式,必须所有人都同意 |
5.2 各种方案总结对比
1. 两阶段提交(2PC) - 像婚礼仪式
-
流程:准备阶段 → 提交/回滚阶段
-
优点:强一致,绝对安全
-
缺点:同步阻塞,性能差,单点故障
-
场景:传统银行、金融系统
2. TCC模式 - 像餐厅预订
-
流程:Try → Confirm/Cancel
-
优点:性能好,资源不长期锁定
-
缺点:实现复杂,业务侵入性强
-
场景:电商交易、支付系统
3. SAGA模式 - 像旅行规划
-
流程:正向操作 → 失败则反向补偿
-
优点:适合长流程业务
-
缺点:补偿操作复杂
-
场景:保险理赔、旅行预订
4. 本地消息表 - 像快递登记
-
流程:业务+消息同事务 → 异步处理
-
优点:简单可靠,业务侵入小
-
缺点:消息有延迟
-
场景:大多数互联网业务
5. 可靠消息队列 - 像邮政系统
-
流程:半消息 → 本地事务 → 确认发送
-
优点:解耦,消息零丢失
-
缺点:需要MQ支持事务
-
场景:需要高可靠的消息传递
5.3 一致性方案对比表
| 方案 | 一致性 | 性能 | 复杂度 | 业务侵入 | 适用场景 |
|---|---|---|---|---|---|
| 2PC | 强一致 | 差 | 中 | 低 | 传统金融系统 |
| TCC | 最终一致 | 好 | 高 | 高 | 电商、互金 |
| SAGA | 最终一致 | 好 | 高 | 中 | 长流程业务 |
| 本地消息表 | 最终一致 | 中 | 中 | 中 | 大多数互联网业务 |
| 可靠消息 | 最终一致 | 好 | 低 | 低 | 高并发场景 |
| 分布式锁 | 强一致 | 中 | 低 | 低 | 配置管理、分布式锁 |
第六章:避坑指南
6.1 常见陷阱与解决方案
陷阱1:分布式锁超时问题
// 错误示例
public void processWithLock() {
String lockKey = "resource_lock";
String requestId = UUID.randomUUID().toString();
try {
// 获取锁,设置30秒超时
boolean locked = redisLock.lock(lockKey, requestId, 30, TimeUnit.SECONDS);
if (!locked) return;
// 业务处理需要60秒
processBusiness(); // 需要60秒
// 问题:锁在第30秒就自动释放了!
// 其他线程在第31秒进入,数据不一致!
} finally {
redisLock.unlock(lockKey, requestId);
}
}
// 正确解决方案:锁续期
public void processWithLockRenewal() {
String lockKey = "resource_lock";
String requestId = UUID.randomUUID().toString();
AtomicBoolean processing = new AtomicBoolean(true);
try {
// 获取锁
boolean locked = redisLock.lock(lockKey, requestId, 30, TimeUnit.SECONDS);
if (!locked) return;
// 启动锁续期线程
Thread renewalThread = new Thread(() -> {
while (processing.get()) {
try {
Thread.sleep(10000); // 每10秒续期一次
redisLock.renew(lockKey, requestId, 30, TimeUnit.SECONDS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
renewalThread.start();
// 执行业务
processBusiness();
} finally {
processing.set(false);
redisLock.unlock(lockKey, requestId);
}
}
陷阱2:消息重复消费
// 错误:可能重复消费
@MQListener(topic = "ORDER_CREATED")
public void handleOrder(OrderMessage message) {
// 直接处理订单
orderService.createOrder(message);
// 如果消息重复,订单会被创建多次!
}
// 正确:保证幂等性
@MQListener(topic = "ORDER_CREATED")
public void handleOrderIdempotent(OrderMessage message) {
// 1. 使用业务唯一标识检查幂等性
if (messageLogService.exists(message.getMessageId())) {
return; // 已处理过
}
// 2. 使用数据库唯一约束
try {
orderService.createOrderWithUniqueConstraint(message);
} catch (DuplicateKeyException e) {
// 唯一约束冲突,说明已处理过
return;
}
// 3. 记录已处理的消息
messageLogService.save(message.getMessageId());
}
陷阱3:补偿操作失败
// 错误:补偿操作可能失败
public void compensateOrder(Order order) {
try {
// 恢复库存
inventoryService.restoreStock(order.getProductId());
// 退款
paymentService.refund(order.getAmount());
} catch (Exception e) {
// 如果这里失败,数据就不一致了!
// 库存恢复了,但款没退
}
}
// 正确:补偿操作要可重试和可监控
public void compensateOrderWithRetry(Order order) {
int maxRetries = 3;
int retryCount = 0;
boolean success = false;
while (retryCount < maxRetries && !success) {
try {
// 使用本地事务保证两个操作原子性
transactionTemplate.execute(status -> {
inventoryService.restoreStock(order.getProductId());
paymentService.refund(order.getAmount());
return null;
});
success = true;
log.info("补偿操作成功: orderId={}", order.getOrderId());
} catch (Exception e) {
retryCount++;
log.warn("补偿操作失败,第{}次重试: orderId={}, error={}",
retryCount, order.getOrderId(), e.getMessage());
if (retryCount >= maxRetries) {
// 达到最大重试次数,告警人工处理
alertService.sendAlert("补偿操作失败需要人工处理",
"orderId=" + order.getOrderId());
} else {
// 指数退避等待
try {
Thread.sleep((long) Math.pow(2, retryCount) * 1000);
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
}
}
}
}
}
6.2 监控与告警最佳实践
@Component
public class MonitoringSystem {
@Autowired
private MeterRegistry meterRegistry;
@Autowired
private AlertService alertService;
// 1. 关键指标监控
@Scheduled(fixedRate = 10000)
public void monitorKeyMetrics() {
// 监控QPS
double qps = getCurrentQPS();
if (qps > 10000) {
alertService.sendAlert("HIGH_QPS", "当前QPS: " + qps);
}
// 监控库存不一致
Map<String, Long> inconsistentStock = checkStockConsistency();
if (!inconsistentStock.isEmpty()) {
alertService.sendAlert("STOCK_INCONSISTENCY",
"不一致的商品: " + inconsistentStock);
}
// 监控消息积压
long backlogCount = mqService.getBacklogCount("ORDER_TOPIC");
if (backlogCount > 1000) {
alertService.sendAlert("MQ_BACKLOG",
"消息积压数量: " + backlogCount);
}
}
// 2. 分布式追踪
public Order createOrderWithTrace(OrderRequest request) {
Span span = tracer.buildSpan("create_order")
.withTag("userId", request.getUserId())
.withTag("productId", request.getProductId())
.start();
try (Scope scope = tracer.activateSpan(span)) {
// 记录开始时间
long startTime = System.currentTimeMillis();
// 执行业务
Order order = orderService.createOrder(request);
// 记录耗时
long duration = System.currentTimeMillis() - startTime;
span.log("create_order_completed, duration=" + duration + "ms");
// 记录指标
meterRegistry.timer("order.create.duration").record(duration, TimeUnit.MILLISECONDS);
return order;
} catch (Exception e) {
span.log(Collections.singletonMap("error", e.getMessage()));
span.setTag("error", true);
throw e;
} finally {
span.finish();
}
}
// 3. 健康检查
@GetMapping("/health")
public HealthResponse healthCheck() {
HealthResponse response = new HealthResponse();
// 检查数据库连接
response.setDatabaseHealthy(checkDatabaseHealth());
// 检查Redis连接
response.setRedisHealthy(checkRedisHealth());
// 检查MQ连接
response.setMqHealthy(checkMQHealth());
// 检查磁盘空间
response.setDiskHealthy(checkDiskSpace());
return response;
}
private double getCurrentQPS() {
// 获取当前QPS
return 0.0;
}
private Map<String, Long> checkStockConsistency() {
// 检查库存一致性
return new HashMap<>();
}
// 健康检查响应类
static class HealthResponse {
private boolean databaseHealthy;
private boolean redisHealthy;
private boolean mqHealthy;
private boolean diskHealthy;
// getters and setters
public boolean isDatabaseHealthy() { return databaseHealthy; }
public void setDatabaseHealthy(boolean databaseHealthy) { this.databaseHealthy = databaseHealthy; }
public boolean isRedisHealthy() { return redisHealthy; }
public void setRedisHealthy(boolean redisHealthy) { this.redisHealthy = redisHealthy; }
public boolean isMqHealthy() { return mqHealthy; }
public void setMqHealthy(boolean mqHealthy) { this.mqHealthy = mqHealthy; }
public boolean isDiskHealthy() { return diskHealthy; }
public void setDiskHealthy(boolean diskHealthy) { this.diskHealthy = diskHealthy; }
}
}
结语:从理论到实践的完整路径
通过这篇长文,我们从最基础的概念开始,逐步深入到分布式一致性的各个层面。让我们总结一下关键收获:
核心原则:
-
没有银弹:每种方案都有优缺点,要根据业务选择
-
先理解业务:技术为业务服务,不是业务为技术服务
-
接受不完美:在一致性、可用性、性能之间找到平衡点
-
监控比功能重要:没有监控的系统就像盲人开车
-
补偿机制必须:任何自动化的东西都可能失败,要有补偿
学习路径建议:
第一阶段(0-6个月):
-
掌握数据库事务和锁
-
理解基本的并发问题
-
学会使用Redis分布式锁
第二阶段(6-18个月):
-
深入理解CAP/BASE理论
-
掌握消息队列的使用
-
实现最终一致性方案
第三阶段(18-36个月):
-
掌握多种分布式事务模式(2PC、TCC、SAGA)
-
能根据业务设计合适的架构
-
建立完整的监控和运维体系
架构师阶段:
-
全局视角,业务驱动
-
技术选型,团队赋能
-
风险控制,前瞻规划
最后的思考题:
假设你要设计以下系统,你会如何选择一致性方案?
-
在线文档协作(如腾讯文档):多人同时编辑,要实时看到对方的修改
-
智能家居控制:手机控制家电,要保证设备状态准确
-
股票交易系统:买卖股票,绝对不能出错
-
社交游戏排名:全球玩家排名,允许短暂不一致
把你的想法写在评论区,我们一起探讨!
如果你在实践过程中遇到问题,或者有更好的想法,欢迎在评论区留言。每一条评论我都会认真阅读和回复!
更多推荐

所有评论(0)