分布式事务问题背景与Saga模式原理
微服务架构中,一个业务操作跨越多个服务节点,每个服务维护独立数据库,传统本地ACID事务无法保证跨服务的数据一致性。分布式事务的核心挑战在于:部分服务执行成功而部分失败时,如何回滚已提交的本地事务。Saga模式将长事务拆分为一系列本地事务,每个本地事务有对应的补偿操作。正向执行时按顺序提交T1、T2、T3,任一步骤失败时反向执行C2、C1补偿已提交的事务。
Saga有两种协调方式:编排式(Orchestration)和协同式(Choreography)。编排式由中央协调器统一调度各服务执行顺序,适合复杂业务流程;协同式通过事件驱动各服务自主响应,适合简单流程。Spring Boot框架下,Seata中间件提供了开箱即用的Saga实现。
Seata Saga模式架构与部署
Seata是阿里巴巴开源的分布式事务解决方案,支持AT、TCC、SAGA、XA四种模式。Saga模式适用于长流程业务场景,如订单创建涉及库存扣减、优惠券核销、积分发放、支付处理等多个步骤。
Seata Server部署(Docker方式):
docker run -d --name seata-server \
-p 8091:8091 \
-p 7091:7091 \
-e SEATA_IP=127.0.0.1 \
-e STORE_MODE=db \
-e STORE_DB_DRIVER_CLASS_NAME=com.mysql.cj.jdbc.Driver \
-e STORE_DB_URL="jdbc:mysql://127.0.0.1:3306/seata?useUnicode=true&characterEncoding=utf8" \
-e STORE_DB_USER=root \
-e STORE_DB_PASSWORD='your_password' \
seataio/seata-server:2.0.0
业务数据库需创建seata全局事务表和分支事务表:
CREATE TABLE IF NOT EXISTS `global_table` (
`xid` varchar(128) NOT NULL,
`transaction_id` bigint(20) DEFAULT NULL,
`status` tinyint(4) NOT NULL,
`application_id` varchar(32) DEFAULT NULL,
`transaction_service_group` varchar(32) DEFAULT NULL,
`transaction_name` varchar(128) DEFAULT NULL,
`timeout` int(11) DEFAULT NULL,
`begin_time` bigint(20) DEFAULT NULL,
`application_data` varchar(2000) DEFAULT NULL,
`gmt_create` datetime DEFAULT NULL,
`gmt_modified` datetime DEFAULT NULL,
PRIMARY KEY (`xid`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
CREATE TABLE IF NOT EXISTS `branch_table` (
`branch_id` bigint(20) NOT NULL,
`xid` varchar(128) NOT NULL,
`transaction_id` bigint(20) DEFAULT NULL,
`resource_group_id` varchar(32) DEFAULT NULL,
`resource_id` varchar(256) DEFAULT NULL,
`branch_type` varchar(8) DEFAULT NULL,
`status` tinyint(4) DEFAULT NULL,
`client_id` varchar(64) DEFAULT NULL,
`application_data` varchar(2000) DEFAULT NULL,
`gmt_create` datetime DEFAULT NULL,
`gmt_modified` datetime DEFAULT NULL,
PRIMARY KEY (`branch_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
Spring Boot集成Seata Saga实战
以电商订单流程为例,定义Saga编排流程。订单服务作为入口,依次调用库存服务、优惠券服务、积分服务。
Maven依赖配置:
<dependency>
<groupId>io.seata</groupId>
<artifactId>seata-spring-boot-starter</artifactId>
<version>2.0.0</version>
</dependency>
<dependency>
<groupId>io.seata</groupId>
<artifactId>seata-saga-engine</artifactId>
<version>2.0.0</version>
</dependency>
application.yml配置:
seata:
enabled: true
application-id: order-service
tx-service-group: my_tx_group
service:
vgroup-mapping:
my_tx_group: default
grouplist:
default: 127.0.0.1:8091
saga:
enabled: true
state-machine:
enable-async: false
trans-operation-timeout: 1800000
Saga状态机JSON定义(定义正向和补偿流程):
{
"Name": "createOrderSaga",
"Comment": "订单创建Saga流程",
"StartState": "deductInventory",
"States": {
"deductInventory": {
"Type": "ServiceTask",
"ServiceName": "inventoryService",
"ServiceMethod": "deduct",
"CompensateState": "compensateInventory",
"Next": "useCoupon"
},
"compensateInventory": {
"Type": "ServiceTask",
"ServiceName": "inventoryService",
"ServiceMethod": "addBack"
},
"useCoupon": {
"Type": "ServiceTask",
"ServiceName": "couponService",
"ServiceMethod": "use",
"CompensateState": "compensateCoupon",
"Next": "addPoints"
},
"compensateCoupon": {
"Type": "ServiceTask",
"ServiceName": "couponService",
"ServiceMethod": "restore"
},
"addPoints": {
"Type": "ServiceTask",
"ServiceName": "pointsService",
"ServiceMethod": "add",
"CompensateState": "compensatePoints",
"Next": "succeed"
},
"compensatePoints": {
"Type": "ServiceTask",
"ServiceName": "pointsService",
"ServiceMethod": "deduct"
},
"succeed": {
"Type": "Succeed"
}
}
}
Java服务实现与触发:
@Service
public class OrderSagaService {
@Autowired
private StateMachineEngine stateMachineEngine;
public String createOrder(OrderDTO order) {
Map<String, Object> params = new HashMap<>();
params.put("orderId", order.getOrderId());
params.put("productId", order.getProductId());
params.put("quantity", order.getQuantity());
params.put("couponId", order.getCouponId());
params.put("userId", order.getUserId());
StateMachineInstance instance = stateMachineEngine.start(
"createOrderSaga",
BusinessType.ORDER.name(),
params
);
if (instance.getStatus() == ExecutionStatus.SU) {
return "订单创建成功: " + order.getOrderId();
} else {
throw new RuntimeException("订单创建失败,已执行补偿: "
+ instance.getException());
}
}
}
@Service
public class InventoryServiceImpl implements InventoryService {
@Autowired
private InventoryMapper inventoryMapper;
@Override
public boolean deduct(Map<String, Object> params) {
String productId = (String) params.get("productId");
Integer quantity = (Integer) params.get("quantity");
int rows = inventoryMapper.deductStock(productId, quantity);
if (rows == 0) {
throw new RuntimeException("库存不足: " + productId);
}
return true;
}
@Override
public boolean addBack(Map<String, Object> params) {
String productId = (String) params.get("productId");
Integer quantity = (Integer) params.get("quantity");
inventoryMapper.addBackStock(productId, quantity);
return true;
}
}
消息中间件在Saga中的可靠性保障
Saga模式中补偿操作的触发依赖协调器与服务间的通信可靠性。引入消息中间件(如RocketMQ)作为服务间通信通道,确保补偿指令不丢失。RocketMQ事务消息机制保证本地事务与消息发送的原子性:
@Service
public class OrderTransactionListener implements TransactionListener {
@Autowired
private OrderService orderService;
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
OrderDTO order = JSON.parseObject(
new String(msg.getBody()), OrderDTO.class);
orderService.createOrder(order);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
String orderId = msg.getKeys();
Order order = orderService.getById(orderId);
if (order != null && order.getStatus() == OrderStatus.CREATED) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.UNKNOW;
}
}
服务治理与幂等性设计
Saga补偿操作必须保证幂等性,因为网络超时可能导致重试。补偿方法通过状态机版本号判断是否已执行,避免重复补偿导致数据不一致。业务中台建设中,建议为每个Saga步骤维护状态表,记录正向和补偿操作的执行状态:
CREATE TABLE saga_step_log (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
saga_id VARCHAR(64) NOT NULL,
step_name VARCHAR(64) NOT NULL,
step_type VARCHAR(16) NOT NULL COMMENT 'FORWARD/COMPENSATION',
status VARCHAR(16) NOT NULL COMMENT 'INIT/RUNNING/SUCCESS/FAIL',
retry_count INT DEFAULT 0,
request_data TEXT,
response_data TEXT,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_saga_step (saga_id, step_name, step_type)
);
补偿方法执行前检查该步骤是否已成功补偿,已成功则直接返回,避免重复执行。API接口规范方面,Saga协调器调用各服务接口统一使用POST方法,请求体携带sagaId和stepName字段,服务端通过这两个字段实现幂等控制。
高并发设计场景下,Saga模式需配合分布式锁防止同一业务流程并发执行。使用Redisson分布式锁在订单维度加锁,确保同一订单不会同时触发多个Saga实例。锁粒度尽量细化到业务实体级别(如orderId),避免全局锁影响吞吐量。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/fen-bu-shi-shi-wu-saga-mo-shi-shi-zhan-springboot-wei-fu-wu/