引言:从奶茶店的故事说起

假设你开了一家网红奶茶店,生意火爆。今天是周六,顾客排起了长队。这时候,你会遇到几个问题:

  1. 只有一个收银员:收银员既要收钱,又要记库存,忙不过来。这就像单体应用

  2. 你增加了收银员:现在有两个收银员,但他们都用一个账本记录库存。这就像多线程并发

  3. 你开了分店:总店和分店共享库存,但沟通靠打电话。这就像分布式系统

  4. 节假日高峰期:两家店同时有100人排队,都来买同一款限量奶茶。这就像高并发

  5. 最头疼的问题:总店卖了最后一杯奶茶,还没来得及告诉分店,分店也卖了这最后一杯。结果超卖了,顾客生气。这就是数据不一致

今天,我就用最通俗的语言,带你彻底理解这个复杂的分布式高并发数据一致性问题。无论你是刚入门的小白,还是想深入理解的老手,这篇文章都会给你清晰的答案。

第一章:基础篇 - 这到底是什么问题?

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的优缺点

  • 优点:强一致性保证

  • 缺点:

    1. 同步阻塞:所有参与者都要等待

    2. 单点故障:协调者挂了,所有参与者卡住

    3. 数据不一致:协调者发送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(分区容错性):部分路段封闭时,其他路段还能通行

  • 一条路修路,其他路还能走

残酷的现实:你只能三选二!

  1. CP系统(一致性+分区容错)

    • 像银行系统:宁可暂时不能转账,也要保证账目绝对正确

    • 网络出问题时,停止服务,保证数据一致

  2. AP系统(可用性+分区容错)

    • 像电商网站:宁可显示旧价格,也要保证能访问

    • 网络出问题时,继续服务,返回可能过期的数据

  3. 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; }
    }
}

结语:从理论到实践的完整路径

通过这篇长文,我们从最基础的概念开始,逐步深入到分布式一致性的各个层面。让我们总结一下关键收获:

核心原则:

  1. 没有银弹:每种方案都有优缺点,要根据业务选择

  2. 先理解业务:技术为业务服务,不是业务为技术服务

  3. 接受不完美:在一致性、可用性、性能之间找到平衡点

  4. 监控比功能重要:没有监控的系统就像盲人开车

  5. 补偿机制必须:任何自动化的东西都可能失败,要有补偿

学习路径建议:

第一阶段(0-6个月)

  • 掌握数据库事务和锁

  • 理解基本的并发问题

  • 学会使用Redis分布式锁

第二阶段(6-18个月)

  • 深入理解CAP/BASE理论

  • 掌握消息队列的使用

  • 实现最终一致性方案

第三阶段(18-36个月)

  • 掌握多种分布式事务模式(2PC、TCC、SAGA)

  • 能根据业务设计合适的架构

  • 建立完整的监控和运维体系

架构师阶段

  • 全局视角,业务驱动

  • 技术选型,团队赋能

  • 风险控制,前瞻规划

最后的思考题:

假设你要设计以下系统,你会如何选择一致性方案?

  1. 在线文档协作(如腾讯文档):多人同时编辑,要实时看到对方的修改

  2. 智能家居控制:手机控制家电,要保证设备状态准确

  3. 股票交易系统:买卖股票,绝对不能出错

  4. 社交游戏排名:全球玩家排名,允许短暂不一致

把你的想法写在评论区,我们一起探讨!

如果你在实践过程中遇到问题,或者有更好的想法,欢迎在评论区留言。每一条评论我都会认真阅读和回复!

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐