📖 数据密集型设计

分布式事务一致性保障

深入探讨分布式事务一致性保障机制

一、分布式事务概述

分布式事务是指涉及多个数据库或服务的事务操作。在分布式系统中,保证事务的ACID特性面临挑战,需要采用特殊的一致性保障机制。

二、分布式事务一致性模型

2.1 一致性模型对比

一致性模型 描述 一致性强度 性能 适用场景
强一致性 所有节点同时看到相同数据 最高 最低 金融交易
顺序一致性 操作按顺序执行 分布式锁
因果一致性 因果相关操作保持顺序 中高 社交网络
最终一致性 最终所有节点数据一致 最高 一般业务

三、两阶段提交(2PC)

3.1 2PC原理

两阶段提交是一种分布式事务协议,分为准备阶段和提交阶段:

sequenceDiagram participant TM as 事务管理器 participant RM1 as 资源管理器1 participant RM2 as 资源管理器2 TM->>RM1: Prepare RM1->>RM1: 执行事务操作 RM1->>RM1: 记录事务日志 RM1-->>TM: Prepared TM->>RM2: Prepare RM2->>RM2: 执行事务操作 RM2->>RM2: 记录事务日志 RM2-->>TM: Prepared alt 所有资源管理器准备成功 TM->>RM1: Commit RM1->>RM1: 提交事务 RM1-->>TM: Committed TM->>RM2: Commit RM2->>RM2: 提交事务 RM2-->>TM: Committed TM-->>TM: 事务完成 else 任一资源管理器准备失败 TM->>RM1: Rollback RM1->>RM1: 回滚事务 RM1-->>TM: Rolledback TM->>RM2: Rollback RM2->>RM2: 回滚事务 RM2-->>TM: Rolledback TM-->>TM: 事务回滚 end

3.2 2PC实现

public class TwoPhaseCommitService
{
    private readonly List<IResourceManager> _resourceManagers;
    
    public async Task<TransactionResult> ExecuteTransactionAsync(
        List<TransactionOperation> operations)
    {
        var transactionId = Guid.NewGuid().ToString();
        
        var prepareResults = await PreparePhaseAsync(transactionId, operations);
        
        if (prepareResults.All(r => r.IsPrepared))
        {
            return await CommitPhaseAsync(transactionId, operations);
        }
        
        return await RollbackPhaseAsync(transactionId, operations);
    }
    
    private async Task<List<PrepareResult>> PreparePhaseAsync(
        string transactionId, 
        List<TransactionOperation> operations)
    {
        var results = new List<PrepareResult>();
        
        foreach (var operation in operations)
        {
            var rm = GetResourceManager(operation.ResourceType);
            
            var result = await rm.PrepareAsync(transactionId, operation);
            results.Add(result);
            
            if (!result.IsPrepared)
                return results;
        }
        
        return results;
    }
    
    private async Task<TransactionResult> CommitPhaseAsync(
        string transactionId, 
        List<TransactionOperation> operations)
    {
        foreach (var operation in operations)
        {
            var rm = GetResourceManager(operation.ResourceType);
            await rm.CommitAsync(transactionId);
        }
        
        return new TransactionResult { Success = true };
    }
    
    private async Task<TransactionResult> RollbackPhaseAsync(
        string transactionId, 
        List<TransactionOperation> operations)
    {
        foreach (var operation in operations)
        {
            var rm = GetResourceManager(operation.ResourceType);
            await rm.RollbackAsync(transactionId);
        }
        
        return new TransactionResult { Success = false };
    }
}

3.3 2PC问题分析

问题 描述 影响 解决方案
同步阻塞 所有参与者等待其他参与者响应 性能低 异步提交
单点故障 事务管理器故障导致无法提交 数据不一致 多副本部署
数据不一致 提交阶段部分成功 数据不一致 补偿事务

四、三阶段提交(3PC)

4.1 3PC原理

三阶段提交在2PC基础上增加了预提交阶段,减少同步阻塞:

sequenceDiagram participant TM as 事务管理器 participant RM1 as 资源管理器1 participant RM2 as 资源管理器2 TM->>RM1: CanCommit RM1-->>TM: Yes TM->>RM2: CanCommit RM2-->>TM: Yes TM->>RM1: PreCommit RM1->>RM1: 执行事务 RM1-->>TM: Prepared TM->>RM2: PreCommit RM2->>RM2: 执行事务 RM2-->>TM: Prepared TM->>RM1: DoCommit RM1->>RM1: 提交事务 RM1-->>TM: Committed TM->>RM2: DoCommit RM2->>RM2: 提交事务 RM2-->>TM: Committed

4.2 3PC与2PC对比

特性 2PC 3PC
阶段数 2 3
同步阻塞
故障处理 复杂 简单
性能

五、Saga模式

5.1 Saga模式原理

Saga模式将分布式事务拆分为一系列本地事务,通过补偿机制保证最终一致性:

flowchart TD A[开始] --> B[本地事务1] B --> C{成功} C -->|是| D[本地事务2] C -->|否| E[回滚] D --> F{成功} F -->|是| G[本地事务3] F -->|否| H[补偿事务2] H --> E G --> I{成功} I -->|是| J[完成] I -->|否| K[补偿事务3] K --> H

5.2 Saga模式实现

public class SagaPatternService
{
    public async Task<SagaResult> ExecuteSagaAsync(Saga saga)
    {
        var executedTransactions = new List<SagaStep>();
        
        try
        {
            foreach (var step in saga.Steps)
            {
                var result = await ExecuteStepAsync(step);
                
                if (!result.Success)
                {
                    await CompensateAsync(executedTransactions);
                    return new SagaResult { Success = false, ErrorMessage = result.ErrorMessage };
                }
                
                executedTransactions.Add(step);
            }
            
            return new SagaResult { Success = true };
        }
        catch (Exception ex)
        {
            await CompensateAsync(executedTransactions);
            return new SagaResult { Success = false, ErrorMessage = ex.Message };
        }
    }
    
    private async Task<StepResult> ExecuteStepAsync(SagaStep step)
    {
        try
        {
            await step.Operation();
            return new StepResult { Success = true };
        }
        catch (Exception ex)
        {
            return new StepResult { Success = false, ErrorMessage = ex.Message };
        }
    }
    
    private async Task CompensateAsync(List<SagaStep> executedTransactions)
    {
        foreach (var step in executedTransactions.AsEnumerable().Reverse())
        {
            try
            {
                if (step.Compensation != null)
                {
                    await step.Compensation();
                }
            }
            catch (Exception ex)
            {
                await _logger.LogErrorAsync($"补偿失败: {ex.Message}");
            }
        }
    }
}

public class Saga
{
    public string SagaId { get; set; }
    public List<SagaStep> Steps { get; set; } = new List<SagaStep>();
}

public class SagaStep
{
    public string StepId { get; set; }
    public Func<Task> Operation { get; set; }
    public Func<Task> Compensation { get; set; }
}

5.3 Order Saga示例

public class OrderSagaService
{
    public async Task<SagaResult> CreateOrderAsync(Order order)
    {
        var saga = new Saga { SagaId = Guid.NewGuid().ToString() };
        
        saga.Steps.Add(new SagaStep
        {
            StepId = "1",
            Operation = async () => await CreateOrderRecordAsync(order),
            Compensation = async () => await CancelOrderAsync(order.Id)
        });
        
        saga.Steps.Add(new SagaStep
        {
            StepId = "2",
            Operation = async () => await ReserveInventoryAsync(order),
            Compensation = async () => await ReleaseInventoryAsync(order)
        });
        
        saga.Steps.Add(new SagaStep
        {
            StepId = "3",
            Operation = async () => await ProcessPaymentAsync(order),
            Compensation = async () => await RefundPaymentAsync(order)
        });
        
        saga.Steps.Add(new SagaStep
        {
            StepId = "4",
            Operation = async () => await NotifyUserAsync(order),
            Compensation = async () => await CancelNotificationAsync(order)
        });
        
        return await _sagaService.ExecuteSagaAsync(saga);
    }
}

六、TCC模式

6.1 TCC模式原理

TCC(Try-Confirm-Cancel)模式将业务操作分为三个阶段:

flowchart TD A[Try阶段] --> B[预留资源] B --> C[检查条件] C -->|成功| D[Confirm阶段] C -->|失败| E[Cancel阶段] D --> F[确认操作] E --> G[释放资源]

6.2 TCC模式实现

public class TccPatternService
{
    public async Task<TccResult> ExecuteTccAsync(TccTransaction transaction)
    {
        var tryResults = await ExecuteTryPhaseAsync(transaction);
        
        if (tryResults.All(r => r.Success))
        {
            return await ExecuteConfirmPhaseAsync(transaction);
        }
        
        return await ExecuteCancelPhaseAsync(transaction, tryResults);
    }
    
    private async Task<List<TccStepResult>> ExecuteTryPhaseAsync(TccTransaction transaction)
    {
        var results = new List<TccStepResult>();
        
        foreach (var participant in transaction.Participants)
        {
            var result = await participant.TryAsync();
            results.Add(result);
            
            if (!result.Success)
                break;
        }
        
        return results;
    }
    
    private async Task<TccResult> ExecuteConfirmPhaseAsync(TccTransaction transaction)
    {
        foreach (var participant in transaction.Participants)
        {
            await participant.ConfirmAsync();
        }
        
        return new TccResult { Success = true };
    }
    
    private async Task<TccResult> ExecuteCancelPhaseAsync(TccTransaction transaction, 
        List<TccStepResult> tryResults)
    {
        for (int i = 0; i < tryResults.Count; i++)
        {
            if (tryResults[i].Success)
            {
                await transaction.Participants[i].CancelAsync();
            }
        }
        
        return new TccResult { Success = false };
    }
}

public interface ITccParticipant
{
    Task<TccStepResult> TryAsync();
    Task ConfirmAsync();
    Task CancelAsync();
}

6.3 TCC参与者实现

public class PaymentParticipant : ITccParticipant
{
    public async Task<TccStepResult> TryAsync()
    {
        var result = await _paymentService.ReservePaymentAsync();
        
        if (result.Success)
        {
            await _transactionLogService.LogAsync("PaymentReserved");
            return new TccStepResult { Success = true };
        }
        
        return new TccStepResult { Success = false };
    }
    
    public async Task ConfirmAsync()
    {
        await _paymentService.ConfirmPaymentAsync();
        await _transactionLogService.LogAsync("PaymentConfirmed");
    }
    
    public async Task CancelAsync()
    {
        await _paymentService.ReleasePaymentAsync();
        await _transactionLogService.LogAsync("PaymentCancelled");
    }
}

七、本地消息表模式

7.1 本地消息表原理

本地消息表模式通过数据库事务保证消息的可靠性:

sequenceDiagram participant App as 应用程序 participant DB as 数据库 participant MessageQueue as 消息队列 participant Service as 下游服务 App->>DB: 本地事务(业务操作+插入消息) DB-->>App: 事务成功 App->>MessageQueue: 发送消息 MessageQueue-->>App: 发送成功 App->>DB: 删除消息记录 DB-->>App: 删除成功 MessageQueue->>Service: 消费消息 Service->>Service: 执行业务逻辑 Service-->>MessageQueue: 确认消费

7.2 本地消息表实现

public class LocalMessageTableService
{
    public async Task SendMessageAsync(string messageType, object messageData)
    {
        using var transaction = await _dbContext.Database.BeginTransactionAsync();
        
        try
        {
            var message = new LocalMessage
            {
                Id = Guid.NewGuid(),
                MessageType = messageType,
                MessageData = JsonSerializer.Serialize(messageData),
                Status = MessageStatus.Pending,
                CreatedAt = DateTime.Now
            };
            
            await _dbContext.LocalMessages.AddAsync(message);
            await _dbContext.SaveChangesAsync();
            
            await transaction.CommitAsync();
            
            await ProcessPendingMessagesAsync();
        }
        catch
        {
            await transaction.RollbackAsync();
            throw;
        }
    }
    
    public async Task ProcessPendingMessagesAsync()
    {
        var pendingMessages = await _dbContext.LocalMessages
            .Where(m => m.Status == MessageStatus.Pending)
            .OrderBy(m => m.CreatedAt)
            .Take(100)
            .ToListAsync();
        
        foreach (var message in pendingMessages)
        {
            await ProcessMessageAsync(message);
        }
    }
    
    private async Task ProcessMessageAsync(LocalMessage message)
    {
        try
        {
            await _messageQueue.SendMessageAsync(message.MessageType, message.MessageData);
            
            message.Status = MessageStatus.Sent;
            message.SentAt = DateTime.Now;
            
            await _dbContext.SaveChangesAsync();
        }
        catch (Exception ex)
        {
            message.RetryCount++;
            
            if (message.RetryCount >= 3)
            {
                message.Status = MessageStatus.Failed;
            }
            
            await _dbContext.SaveChangesAsync();
        }
    }
    
    public async Task CleanupSentMessagesAsync()
    {
        var cutoffDate = DateTime.Now.AddDays(-7);
        
        var sentMessages = await _dbContext.LocalMessages
            .Where(m => m.Status == MessageStatus.Sent && m.SentAt < cutoffDate)
            .ToListAsync();
        
        _dbContext.LocalMessages.RemoveRange(sentMessages);
        await _dbContext.SaveChangesAsync();
    }
}

八、最终一致性实现

8.1 最终一致性保障机制

flowchart TD A[数据变更] --> B[记录事件] B --> C[发布事件] C --> D[消息队列] D --> E[消费者1] D --> F[消费者2] E --> G[处理数据] F --> H[处理数据] G --> I{成功} H --> J{成功} I -->|是| K[确认] I -->|否| L[重试] J -->|是| M[确认] J -->|否| N[重试] L --> G N --> H

8.2 最终一致性监控

public class EventualConsistencyMonitor
{
    public async Task<ConsistencyMetrics> GetMetricsAsync()
    {
        var pendingEvents = await _eventStore.CountPendingEventsAsync();
        var failedEvents = await _eventStore.CountFailedEventsAsync();
        var averageProcessingTime = await _eventStore.GetAverageProcessingTimeAsync();
        
        return new ConsistencyMetrics
        {
            PendingEvents = pendingEvents,
            FailedEvents = failedEvents,
            AverageProcessingTimeMs = averageProcessingTime,
            IsConsistent = pendingEvents == 0 && failedEvents == 0
        };
    }
    
    public async Task MonitorAsync()
    {
        var metrics = await GetMetricsAsync();
        
        if (metrics.PendingEvents > 1000)
        {
            await _alertService.SendAlert("待处理事件过多", 
                $"待处理事件: {metrics.PendingEvents}");
        }
        
        if (metrics.FailedEvents > 100)
        {
            await _alertService.SendAlert("失败事件过多", 
                $"失败事件: {metrics.FailedEvents}");
        }
        
        if (metrics.AverageProcessingTimeMs > 1000)
        {
            await _alertService.SendAlert("事件处理延迟过高", 
                $"平均处理时间: {metrics.AverageProcessingTimeMs}ms");
        }
    }
    
    public async Task RepairFailedEventsAsync()
    {
        var failedEvents = await _eventStore.GetFailedEventsAsync();
        
        foreach (var failedEvent in failedEvents)
        {
            await _eventProcessor.ProcessEventAsync(failedEvent);
        }
    }
}

九、分布式事务最佳实践

9.1 避免分布式事务

通过设计优化,尽量避免需要分布式事务的场景。

9.2 选择合适的一致性模型

根据业务需求选择强一致性或最终一致性。

9.3 使用补偿机制

实现可靠的补偿事务,确保数据最终一致。

9.4 监控事务状态

实时监控分布式事务的状态,及时发现问题。

9.5 测试故障场景

测试各种故障场景,验证事务恢复能力。

十、总结

分布式事务一致性保障是分布式系统设计中的核心挑战。2PC和3PC提供强一致性但性能较低,Saga和TCC模式通过补偿机制实现最终一致性,本地消息表模式利用数据库事务保证消息可靠性。选择合适的一致性模型和事务模式,需要在一致性、性能和复杂度之间取得平衡。