一、分布式事务概述
分布式事务是指涉及多个数据库或服务的事务操作。在分布式系统中,保证事务的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模式通过补偿机制实现最终一致性,本地消息表模式利用数据库事务保证消息可靠性。选择合适的一致性模型和事务模式,需要在一致性、性能和复杂度之间取得平衡。