分布式事务中的TCC与SAGA模式:补偿机制的工程实现与幂等性保证
一、当转账服务在网络超时后出现重复扣款:分布式事务的最终一致性困境
分布式支付系统中,一次跨服务的转账涉及三个操作:账户A扣款 → 账户B入账 → 记录流水。第二步的网络超时导致调用方重试——结果账户B被入账两次。根因在于:标准的ACID事务无法跨越微服务边界,而简单重试机制缺乏幂等性保证。
TCC(Try-Confirm-Cancel)和SAGA是两种主流方案。TCC通过两阶段资源预留提供原子性,适合金融场景。SAGA通过异步补偿链处理长时间运行的业务流程,适合订单处理。两者的共同核心是:每一个操作必须有对应的补偿操作,且补偿操作必须幂等。
二、TCC与SAGA的模式对比
sequenceDiagram
participant C as Coordinator
participant S1 as Service A (账户)
participant S2 as Service B (账户)
Note over C,S2: === TCC模式 ===
C->>S1: Try: 冻结100元
S1–>>C: OK (资源预留)
C->>S2: Try: 预入账100元
S2–>>C: OK
C->>S1: Confirm: 实际扣款
S1–>>C: OK
C->>S2: Confirm: 实际入账
S2–>>C: OK
Note over C,S2: === SAGA模式(订单处理) ===
C->>S1: 创建订单
S1–>>C: OK
C->>S2: 扣减库存
S2–>>C: OK
C->>S3: 创建物流单
S3–>>C: FAIL
Note over C: 触发补偿链
C->>S2: 补偿: 恢复库存
S2–>>C: OK
C->>S1: 补偿: 取消订单
S1–>>C: OK
TCC与SAGA的核心区别:
- TCC:资源先预留(Try),再确认(Confirm/Try不改变数据可见性)
- SAGA:操作立即生效,失败时执行反向补偿操作
- TCC的Try阶段需要业务语义支持——不是所有操作都能"预留"
三、TCC事务的工程实现
use std::collections::HashMap;
use std::sync::Arc;
use serde::{Serialize, Deserialize};
/// 全局事务ID:幂等性保证的核心
/// 格式:{时间戳}-{服务ID}-{自增序号}
/// 通过单调递增和唯一性检测防止重复执行
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
struct TransactionId(String);
impl TransactionId {
fn new(service_id: &str, seq: u64) -> Self {
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis();
Self(format!("{}-{}-{}", ts, service_id, seq))
}
}
/// TCC事务状态
#[derive(Debug, Clone, PartialEq, Eq)]
enum TccPhase {
/// Try阶段:资源已预留,等待Confirm或Cancel
Trying,
/// Confirm阶段:业务操作已最终确认
Confirmed,
/// Cancel阶段:资源已释放,事务回滚
Cancelled,
/// 超时未确认:需要定时任务处理
Timeout,
}
/// Try操作的结果:资源预留信息
#[derive(Debug, Clone, Serialize, Deserialize)]
struct TryResult {
/// 预留的资源ID:用于后续Confirm/Cancel
reservation_id: String,
/// 预留的金额/数量
amount: i64,
/// 预留过期时间:超时后自动Cancel
expire_at: chrono::DateTime<chrono::Utc>,
}
/// TCC事务协调器
struct TccCoordinator {
/// 事务日志:持久化所有TCC事务状态
/// 用于故障恢复和幂等性判断
transaction_log: Arc<RwLock<HashMap<TransactionId, TccTransaction>>>,
/// 已处理的操作记录:幂等性保证
processed_ops: Arc<RwLock<HashMap<String, OperationResult>>>,
}
#[derive(Debug, Clone)]
struct TccTransaction {
txn_id: TransactionId,
phase: TccPhase,
/// 各参与方的Try结果
try_results: Vec<(String, TryResult)>,
/// 创建时间
created_at: chrono::DateTime<chrono::Utc>,
/// 最后更新时间
updated_at: chrono::DateTime<chrono::Utc>,
}
#[derive(Debug, Clone)]
struct OperationResult {
/// 操作是否成功
success: bool,
/// 操作结果数据
data: Option<Vec<u8>>,
/// 执行时间
executed_at: chrono::DateTime<chrono::Utc>,
}
/// TCC事务:账户转账示例
struct AccountTccService {
coordinator: Arc<TccCoordinator>,
}
impl AccountTccService {
/// Try阶段:冻结转出方资金
async fn try_debit(
&self,
txn_id: &TransactionId,
account_id: &str,
amount: i64,
ttl_seconds: u32,
) -> Result<TryResult, TccError> {
// 幂等性检查:如果此操作已执行过,返回缓存结果
let idempotent_key = format!("try_debit:{account_id}:{txn_id.0}");
if let Some(result) = self.get_idempotent_result(&idempotent_key).await {
return match result.success {
true => {
let try_result: TryResult = bincode::deserialize(
&result.data.unwrap()
)?;
Ok(try_result)
}
false => Err(TccError::AlreadyFailed),
};
}
// 实际冻结操作
let reservation = format!("resv_{}_{}", account_id, txn_id.0);
let now = chrono::Utc::now();
// 调用账户服务冻结资金
let success = self.freeze_balance(account_id, amount, &reservation).await?;
if !success {
// 记录失败状态(幂等性保证)
self.record_operation(&idempotent_key, OperationResult {
success: false,
data: None,
executed_at: now,
}).await;
return Err(TccError::InsufficientBalance);
}
let try_result = TryResult {
reservation_id: reservation.clone(),
amount,
expire_at: now + chrono::Duration::seconds(ttl_seconds as i64),
};
// 记录成功状态
let data = bincode::serialize(&try_result)?;
self.record_operation(&idempotent_key, OperationResult {
success: true,
data: Some(data),
executed_at: now,
}).await;
// 记录到事务日志
self.coordinator.log_try(txn_id, "account_a".to_string(), try_result.clone()).await;
Ok(try_result)
}
/// Confirm阶段:实际扣款
async fn confirm_debit(
&self,
txn_id: &TransactionId,
reservation_id: &str,
) -> Result<(), TccError> {
let idempotent_key = format!("confirm_debit:{reservation_id}");
if let Some(result) = self.get_idempotent_result(&idempotent_key).await {
return if result.success { Ok(()) } else { Err(TccError::ConfirmFailed) };
}
// 实际扣款:将冻结资金转为正式扣款
self.commit_freeze(reservation_id).await?;
self.record_operation(&idempotent_key, OperationResult {
success: true,
data: None,
executed_at: chrono::Utc::now(),
}).await;
self.coordinator.update_phase(txn_id, TccPhase::Confirmed).await;
Ok(())
}
/// Cancel阶段:释放冻结资金
async fn cancel_debit(
&self,
txn_id: &TransactionId,
reservation_id: &str,
) -> Result<(), TccError> {
let idempotent_key = format!("cancel_debit:{reservation_id}");
if let Some(result) = self.get_idempotent_result(&idempotent_key).await {
return if result.success { Ok(()) } else { Err(TccError::CancelFailed) };
}
// 释放冻结资金
self.unfreeze_balance(reservation_id).await?;
self.record_operation(&idempotent_key, OperationResult {
success: true,
data: None,
executed_at: chrono::Utc::now(),
}).await;
self.coordinator.update_phase(txn_id, TccPhase::Cancelled).await;
Ok(())
}
/// 幂等性结果获取
async fn get_idempotent_result(&self, key: &str) -> Option<OperationResult> {
self.coordinator.processed_ops
.read()
.await
.get(key)
.cloned()
}
/// 记录操作结果(幂等性保证的核心)
async fn record_operation(&self, key: &str, result: OperationResult) {
self.coordinator.processed_ops
.write()
.await
.insert(key.to_string(), result);
}
// 简化的账户操作接口
async fn freeze_balance(&self, account: &str, amount: i64, reservation: &str) -> Result<bool, TccError> {
Ok(true)
}
async fn commit_freeze(&self, reservation: &str) -> Result<(), TccError> {
Ok(())
}
async fn unfreeze_balance(&self, reservation: &str) -> Result<(), TccError> {
Ok(())
}
}
impl TccCoordinator {
async fn log_try(&self, txn_id: &TransactionId, service: String, result: TryResult) {
let mut log = self.transaction_log.write().await;
let entry = log.entry(txn_id.clone()).or_insert_with(|| TccTransaction {
txn_id: txn_id.clone(),
phase: TccPhase::Trying,
try_results: Vec::new(),
created_at: chrono::Utc::now(),
updated_at: chrono::Utc::now(),
});
entry.try_results.push((service, result));
}
async fn update_phase(&self, txn_id: &TransactionId, phase: TccPhase) {
if let Some(entry) = self.transaction_log.write().await.get_mut(txn_id) {
entry.phase = phase;
entry.updated_at = chrono::Utc::now();
}
}
}
/// SAGA编排器:顺序执行+失败补偿
struct SagaOrchestrator {
/// SAGA步骤定义
steps: Vec<SagaStep>,
/// 已执行步骤的状态:用于补偿
executed_steps: Vec<ExecutedStep>,
}
#[derive(Debug)]
struct SagaStep {
name: String,
/// 正向操作
action: Box<dyn Fn() -> SagaResult + Send + Sync>,
/// 补偿操作(幂等)
compensate: Box<dyn Fn() -> SagaResult + Send + Sync>,
}
#[derive(Debug, Clone)]
struct ExecutedStep {
name: String,
/// 步骤执行结果
result: Option<SagaResult>,
/// 幂等键:用于补偿操作的幂等性
idempotent_key: String,
}
type SagaResult = Result<Vec<u8>, String>;
impl SagaOrchestrator {
/// 执行SAGA:步骤失败后执行补偿链
async fn execute(&mut self) -> Result<(), SagaError> {
for (i, step) in self.steps.iter().enumerate() {
let result = (step.action)();
match result {
Ok(data) => {
self.executed_steps.push(ExecutedStep {
name: step.name.clone(),
result: Some(Ok(data)),
idempotent_key: format!("saga_step:{step_name}:{i}",
step_name = step.name),
});
}
Err(e) => {
// 步骤失败:执行补偿链
tracing::error!(
step = %step.name,
error = %e,
"SAGA step failed, starting compensation"
);
// 反向补偿:从最近执行的步骤开始
self.compensate().await?;
return Err(SagaError::StepFailed {
step: step.name.clone(),
error: e,
});
}
}
}
Ok(())
}
/// 补偿链:反向执行已成功步骤的补偿操作
async fn compensate(&mut self) -> Result<(), SagaError> {
let mut compensation_errors = Vec::new();
while let Some(step) = self.executed_steps.pop() {
let step_def = self.steps.iter()
.find(|s| s.name == step.name)
.ok_or(SagaError::CompensationStepNotFound)?;
let result = (step_def.compensate)();
if let Err(e) = result {
// 补偿失败:记录但继续执行——避免补偿链中断
// 未补偿的操作需要人工介入
compensation_errors.push(format!("{}: {}", step.name, e));
tracing::error!(
step = %step.name,
error = %e,
"SAGA compensation failed, requires manual intervention"
);
}
}
if !compensation_errors.is_empty() {
return Err(SagaError::CompensationFailed(compensation_errors));
}
Ok(())
}
}
#[derive(Debug, thiserror::Error)]
enum TccError {
#[error("Insufficient balance")]
InsufficientBalance,
#[error("TCC operation already failed")]
AlreadyFailed,
#[error("Confirm phase failed")]
ConfirmFailed,
#[error("Cancel phase failed")]
CancelFailed,
#[error("Serialization error: {0}")]
Serialization(#[from] bincode::Error),
}
#[derive(Debug, thiserror::Error)]
enum SagaError {
#[error("SAGA step '{step}' failed: {error}")]
StepFailed { step: String, error: String },
#[error("Compensation step not found")]
CompensationStepNotFound,
#[error("Compensation failed: {0:?}")]
CompensationFailed(Vec<String>),
}
use tokio::sync::RwLock;
核心设计:
- 幂等性通过processed_ops映射实现:每个操作有唯一idempotent_key
- TCC的幂等性在Try/Confirm/Cancel每个阶段独立保证
- SAGA补偿链反向执行:最近成功的步骤先补偿
- 补偿失败不中断链:记录所有失败供人工处理
TCC的空回滚与悬挂检测——幂等性之外的深层问题。空回滚(Empty Rollback)指Cancel操作被调用时,对应的Try操作尚未执行或未完成。例如网络延迟导致Cancel先于Try到达——此时Cancel不能拒绝,必须返回成功(因为调用方认为Cancel是幂等的),但这可能留下"幽灵预留":后续Try到达时,资源虽然被Cancel了,但Try的幂等键未记录,导致Try再次预留成功。解决方案是在Try阶段检查Cancel标记:如果该事务ID已被Cancel过,Try应返回失败。悬挂问题则相反:Try成功后,Confirm/Cancel迟迟不到,资源处于"冻结"状态。需要定时任务扫描超时的Try记录,自动触发Cancel释放资源。另一个工程实践是"SAGA补偿的幂等性粒度"——补偿操作应与正向操作共用幂等键(如saga_step:{step_name}:{txn_id}),但需要区分"正向是否已执行"。如果正向操作网络超时但对端实际已执行,调用方重试正向操作时幂等键命中返回成功——但补偿操作也必须能正确撤销该正向操作的结果。这就要求正向操作的结果中携带足够的信息供补偿操作使用(如操作前后的状态快照),而不是简单依赖"正向未执行→补偿跳过"的假设。这种"操作级状态机"是TCC/SAGA正确性的核心——每个事务ID对应一个状态机,幂等键是状态机的唯一标识,操作结果是状态的快照。
四、TCC与SAGA的适用边界
TCC适用场景:
- 金融交易:需原子性的资源操作(扣款/入账)
- 资源预留:库存、座位、配额等
- 短事务:Try后应在秒级内Confirm/Cancel
SAGA适用场景:
- 长时间运行的业务流程:订单处理、审批流
- 异步解耦的微服务间协调
- 无法提供Try语义的操作(如发送通知)
TCC的局限:
- Try/Confirm/Cancel需要业务代码支持——侵入性强
- 资源冻结期间影响并发操作
- 空回滚问题:Cancel调用时Try尚未执行
幂等性的关键设计:
- idempotent_key必须全局唯一:{操作类型}:{资源ID}:{事务ID}
- 幂等结果需持久化:进程重启后仍能识别已执行操作
- 幂等结果的TTL需大于事务最大执行时间






