Spring Boot微服务架构实战:我是如何用消息中间件解决分布式事务难题的

从一个炸了的订单系统说起

去年我们团队接手了一个电商平台的微服务重构项目,把原来一个单体Spring Boot应用拆成订单服务、库存服务、支付服务、物流服务四五个独立模块。拆完之后功能是解耦了,但一个最头疼的问题冒出来了——下单流程跨了三个服务,怎么保证数据一致性?

我记得上线第一周,凌晨两点被电话叫醒,用户反馈下单成功了但库存没扣。查日志发现,订单服务调用库存服务的HTTP接口超时了,订单写入了但库存没减,出现了超卖。那天晚上我盯着屏幕上两条对不上的数据,意识到分布式事务不是拆完服务就能自然消失的。

这篇文章我就把这一年多踩过的坑、最终落地的方案原原本本讲一遍,包括从XA两阶段提交到RocketMQ事务消息的演进过程、Spring Boot集成RocketMQ的实战代码、幂等库存扣减的Redis方案、用Go写高并发网关的经历,以及Sentinel+Nacos做服务治理的配置细节。

先试了XA两阶段提交,然后放弃了

团队最早讨论的方案是XA两阶段提交。说实话,XA在理论上是最优雅的——全局事务协调者统一管理所有参与者的提交和回滚,要么全成功要么全失败,强一致性保证。我们用了Atomikos作为Spring Boot的XA事务管理器,配置方式大概是这样的:

@Configuration
public class XADataSourceConfig {

    @Bean(name = "orderDataSource")
    @Primary
    public DataSource orderDataSource() {
        AtomikosDataSourceBean ds = new AtomikosDataSourceBean();
        ds.setXaDataSourceClassName("com.mysql.cj.jdbc.MysqlXADataSource");
        ds.setUniqueResourceName("orderDB");
        ds.setXaProperties(xaProps("jdbc:mysql://order-db:3306/order_db"));
        return ds;
    }

    @Bean(name = "inventoryDataSource")
    public DataSource inventoryDataSource() {
        AtomikosDataSourceBean ds = new AtomikosDataSourceBean();
        ds.setXaDataSourceClassName("com.mysql.cj.jdbc.MysqlXADataSource");
        ds.setUniqueResourceName("inventoryDB");
        ds.setXaProperties(xaProps("jdbc:mysql://inventory-db:3306/inventory_db"));
        return ds;
    }

    private Properties xaProps(String url) {
        Properties p = new Properties();
        p.setProperty("url", url);
        p.setProperty("user", "root");
        p.setProperty("password", "secret");
        return p;
    }
}

跑起来之后问题来了。第一,性能惨不忍睹。XA要在第一阶段锁住所有参与者的资源,第二阶段才释放,我们压测时下单TPS从单体的800掉到了120。第二,库存服务和支付服务用的是不同的数据库实例,网络抖动会导致事务协调者超时,然后整个事务卡在中间状态。第三,一旦协调者本身挂了,部分参与者收不到commit或rollback指令,就得人工介入查日志恢复数据。

更让我崩溃的是,我们还有一个支付服务对接的是第三方支付渠道,根本不支持XA协议。这意味着XA的强一致性覆盖不了全链路,那用它还有什么意义?开了几次技术评审会后,我拍板放弃了XA,转向最终一致性方案。

RocketMQ事务消息:最终一致性的最优解

在对比了本地消息表、TCC、Saga几个方案后,我选择了RocketMQ的事务消息。原因很简单:事务消息是RocketMQ原生支持的能力,不需要额外的协调框架,对业务代码侵入最小,而且天然保证本地事务和消息发送的原子性。

事务消息的核心流程我画了个脑图反复给团队讲:先执行本地事务(比如创建订单),然后根据本地事务结果决定是否提交消息。如果本地事务成功,消息投递给消费者触发下游操作(比如扣库存);如果本地事务失败,消息被回滚丢弃。万一本地事务执行完但没来得及返回状态,RocketMQ会回查本地事务状态,保证不会丢消息也不会发错消息。

下面是我实际项目中Spring Boot订单服务的核心代码:

@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    public CreateOrderResult createOrder(OrderRequest request) {
        // 1. 本地创建订单(状态为待处理)
        Order order = new Order();
        order.setOrderNo(generateOrderNo());
        order.setUserId(request.getUserId());
        order.setProductId(request.getProductId());
        order.setQuantity(request.getQuantity());
        order.setStatus("CREATED");
        order.setCreateTime(LocalDateTime.now());
        orderMapper.insert(order);

        // 2. 发送事务消息,触发库存扣减
        String topic = "ORDER_CREATE_TOPIC";
        String tag = "inventory_deduct";
        String destination = topic + ":" + tag;

        OrderMessage msg = new OrderMessage();
        msg.setOrderNo(order.getOrderNo());
        msg.setProductId(order.getProductId());
        msg.setQuantity(order.getQuantity());

        TransactionMQProducer producer = (TransactionMQProducer) rocketMQTemplate.getProducer();
        SendResult sendResult = rocketMQTemplate.sendMessageInTransaction(
            destination,
            MessageBuilder.withPayload(msg)
                .setHeader("orderNo", order.getOrderNo())
                .build(),
            order  // 传递本地参数给事务监听器
        );

        if (sendResult.getSendStatus() != SendStatus.SEND_OK) {
            log.error("事务消息发送失败, orderNo={}", order.getOrderNo());
            throw new BusinessException("下单失败,请重试");
        }

        return new CreateOrderResult(order.getOrderNo(), "处理中");
    }
}

关键点在于sendMessageInTransaction这个方法,它不是普通的消息发送,而是跟一个事务监听器绑定的。RocketMQ会先发一个半消息(Half Message),然后执行本地事务,再根据结果提交或回滚。

@RocketMQTransactionListener的回查机制

事务监听器是整个方案的灵魂。我第一次写的时候漏掉了回查逻辑,结果压测时发现偶尔有订单创建成功但库存消息没发出去,排查才知道是本地事务超时导致RocketMQ没收到确认,回查时我又没正确返回状态,消息就被默认回滚了。

@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQListener {

    @Autowired
    private OrderMapper orderMapper;

    /**
     * 执行本地事务
     */
    @Override
    public void onMessage(String message) {
        // 不在这里做具体逻辑,由executeLocalTransaction处理
    }
}

@Component
@RocketMQTransactionListener
public class OrderTransactionListenerImpl implements RocketMQTransactionListener {

    @Autowired
    private OrderMapper orderMapper;

    /**
     * 执行本地事务 - 在半消息发送成功后自动回调
     */
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        Order order = (Order) arg;
        try {
            // 这里本地事务已经在Service层执行过了
            // 只需要返回COMMIT即可
            return RocketMQLocalTransactionState.COMMIT;
        } catch (Exception e) {
            log.error("本地事务执行异常", e);
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }

    /**
     * 回查本地事务状态
     * 当RocketMQ没有收到commit/rollback确认时,会定期回调此方法
     */
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
        String orderNo = (String) msg.getHeaders().get("orderNo");
        Order order = orderMapper.selectByOrderNo(orderNo);

        if (order == null) {
            // 订单不存在,说明本地事务回滚了
            log.warn("回查: 订单不存在, orderNo={}, 回滚消息", orderNo);
            return RocketMQLocalTransactionState.ROLLBACK;
        }

        if ("CREATED".equals(order.getStatus()) || "PAID".equals(order.getStatus())) {
            // 订单存在且状态有效,提交消息
            log.info("回查: 订单状态有效, orderNo={}, 提交消息", orderNo);
            return RocketMQLocalTransactionState.COMMIT;
        }

        if ("CANCELLED".equals(order.getStatus())) {
            // 订单已取消,回滚消息
            log.info("回查: 订单已取消, orderNo={}, 回滚消息", orderNo);
            return RocketMQLocalTransactionState.ROLLBACK;
        }

        // 未知状态,返回UNKNOW让RocketMQ稍后再查
        return RocketMQLocalTransactionState.UNKNOWN;
    }
}

我对回查逻辑有两个要求:第一,查数据库的次数要尽量少,后面我会讲怎么用Redis缓存减少回查压力;第二,UNKNOW状态不能滥用,必须保证最终能给出确定结果,不然消息就永远悬着了。我给回查设了最大3次的限制,超过3次统一返回ROLLBACK,然后靠对账脚本补偿。

幂等库存扣减:Redis去重+数据库乐观锁

消息消费端最怕的就是重复消费。RocketMQ虽然保证消息至少投递一次,但网络问题可能导致同一消息被投递两次甚至更多次。如果库存扣减不是幂等的,扣两次就出大问题了。

我设计的幂等方案分两层:第一层用Redis做消息去重,第二层用数据库乐观锁兜底。

@Component
@RocketMQMessageListener(
    topic = "ORDER_CREATE_TOPIC",
    selectorExpression = "inventory_deduct",
    consumerGroup = "inventory-consumer-group"
)
public class InventoryDeductConsumer implements RocketMQListener {

    @Autowired
    private StringRedisTemplate redisTemplate;

    @Autowired
    private InventoryMapper inventoryMapper;

    @Override
    public void onMessage(OrderMessage msg) {
        String msgId = msg.getOrderNo() + "_inventory_deduct";
        String redisKey = "deduct_msg_id:" + msgId;

        // 第一层:Redis去重,利用SETNX原子操作
        Boolean isNew = redisTemplate.opsForValue()
            .setIfAbsent(redisKey, "1", Duration.ofHours(24));

        if (Boolean.FALSE.equals(isNew)) {
            log.info("重复消息已处理, msgId={}", msgId);
            return;
        }

        try {
            // 第二层:数据库乐观锁扣减库存
            int updated = inventoryMapper.deductStockWithVersion(
                msg.getProductId(),
                msg.getQuantity()
            );

            if (updated == 0) {
                // 扣减失败,可能是库存不足或版本号冲突
                // 删除Redis标记,让消息可以重试
                redisTemplate.delete(redisKey);
                throw new RuntimeException("库存扣减失败, productId="
                    + msg.getProductId());
            }

            log.info("库存扣减成功, orderNo={}, productId={}, qty={}",
                msg.getOrderNo(), msg.getProductId(), msg.getQuantity());

        } catch (Exception e) {
            // 处理失败,删除Redis标记让消息重试
            redisTemplate.delete(redisKey);
            throw e;  // 抛异常触发RocketMQ重试
        }
    }
}

对应的Mapper SQL是这样的:

@Update("UPDATE inventory SET stock = stock - #{quantity}, "
    + "version = version + 1 "
    + "WHERE product_id = #{productId} "
    + "AND stock >= #{quantity} "
    + "AND version = #{version}")
int deductStockWithVersion(
    @Param("productId") Long productId,
    @Param("quantity") Integer quantity);

这里有个细节我折腾了好久:Redis的SETNX和数据库操作之间如果Redis标记设置了但数据库操作失败了,必须删除Redis标记,否则消息永远不会被重新消费。最初我忘了加异常处理里的redisTemplate.delete,导致偶发的库存扣减失败消息被永久跳过。那次的超卖事故让我彻底理解了”幂等不是去重,去重只是幂等的手段之一”这句话。

Go语言API网关:扛住高并发的第一道关

微服务拆分后,我们需要一个统一的入口做路由转发、限流、鉴权。我对比了Spring Cloud Gateway和Kong,但考虑到我们的峰值QPS能到5万,Java网关的GC停顿在高负载下会拖累延迟,最终决定用Go自己写一个轻量网关。

Go的goroutine模型太适合做这种IO密集型网关了。下面是核心的反向代理和限流逻辑:

package gateway

import (
    "net/http"
    "net/http/httputil"
    "net/url"
    "sync"
    "time"

    "github.com/gin-gonic/gin"
)

// 服务路由表,从Nacos动态获取
var (
    routeTable   = make(map[string]string)
    routeTableMu sync.RWMutex
)

func InitRouter() *gin.Engine {
    r := gin.New()
    r.Use(RateLimitMiddleware(100, 200)) // 每秒100令牌,桶容量200
    r.Use(AuthMiddleware())

    // 通用代理转发
    r.Any("/:service/*path", ProxyHandler)

    return r
}

func ProxyHandler(c *gin.Context) {
    serviceName := c.Param("service")
    path := c.Param("path")

    routeTableMu.RLock()
    target, ok := routeTable[serviceName]
    routeTableMu.RUnlock()

    if !ok {
        c.JSON(http.StatusBadGateway, gin.H{
            "code":    502,
            "message": "service not found: " + serviceName,
        })
        return
    }

    targetURL, _ := url.Parse(target)
    proxy := httputil.NewSingleHostReverseProxy(targetURL)

    // 透传请求头,注入追踪ID
    c.Request.Header.Set("X-Request-Id", generateTraceID())
    c.Request.URL.Path = path

    proxy.ServeHTTP(c.Writer, c.Request)
}

// 令牌桶限流中间件
func RateLimitMiddleware(rate, capacity int) gin.HandlerFunc {
    limiter := NewTokenBucketLimiter(rate, capacity)
    return func(c *gin.Context) {
        if !limiter.Allow() {
            c.JSON(http.StatusTooManyRequests, gin.H{
                "code":    429,
                "message": "too many requests",
            })
            c.Abort()
            return
        }
        c.Next()
    }
}

// Nacos路由表热更新
func WatchNacosRoutes(nacosClient *NacosClient) {
    go func() {
        for {
            services := nacosClient.GetAllServices("DEFAULT_GROUP")
            routeTableMu.Lock()
            for _, svc := range services {
                routeTable[svc.Name] = "http://" + svc.Host + ":" + svc.Port
            }
            routeTableMu.Unlock()
            time.Sleep(5 * time.Second)
        }
    }()
}

这个Go网关部署了两台实例,经过压测单实例能稳定处理3万QPS,P99延迟在12ms以内。比之前用Spring Cloud Gateway时P99降低了60%。Go的内存占用也更友好,两个实例总共才用了不到200MB内存,之前Java网关动不动就吃掉1.5GB。

当然Go网关也不是完美的,它的缺点是生态不如Java丰富,比如JWT校验、请求签名这些能力要自己实现或找第三方库。但作为纯粹的反向代理和流量管控节点,Go的性能优势是压倒性的。

Sentinel + Nacos:微服务治理的左膀右臂

服务拆多了之后,故障隔离变成头号问题。有一次库存服务因为一个慢SQL卡住了,导致调用它的订单服务线程池耗尽,订单服务又拖垮了支付服务,雪崩效应整个平台不可用了半个小时。

我用了两个工具解决这个问题:Sentinel做流量防护,Nacos做配置中心和服务发现。

Sentinel的接入非常简单,Spring Boot项目加个依赖就能用:

<!-- pom.xml -->
<dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-sentinel</artifactId>
</dependency>
<dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-nacos-discovery</artifactId>
</dependency>

然后是application.yml的关键配置:

spring:
  cloud:
    sentinel:
      transport:
        dashboard: sentinel-dashboard:8080
        port: 8719
      datasource:
        # 从Nacos拉取限流规则,实现动态配置
        flow:
          nacos:
            server-addr: nacos:8848
            namespace: production
            group-id: SENTINEL_GROUP
            data-id: order-service-flow-rules
            rule-type: flow
        degrade:
          nacos:
            server-addr: nacos:8848
            namespace: production
            group-id: SENTINEL_GROUP
            data-id: order-service-degrade-rules
            rule-type: degrade

    nacos:
      discovery:
        server-addr: nacos:8848
        namespace: production

Sentinel的熔断降级配置我直接推到Nacos上,这样不用重启服务就能动态调整规则。我们在Nacos上配置的降级规则长这样:

[
    {
        "resource": "InventoryService:deductStock",
        "grade": 1,
        "count": 0.5,
        "timeWindow": 30,
        "minRequestAmount": 10,
        "statIntervalMs": 5000
    },
    {
        "resource": "PaymentService:pay",
        "grade": 1,
        "count": 0.3,
        "timeWindow": 20,
        "minRequestAmount": 5,
        "statIntervalMs": 5000
    }
]

这段配置的意思是:当库存扣减接口的错误率超过50%且请求量大于10时,自动熔断30秒;支付接口错误率超过30%且请求量大于5时,熔断20秒。熔断期间请求直接走降级逻辑,不会打到下游服务。

降级逻辑我用了@SentinelResource注解:

@Service
public class OrderServiceImpl implements OrderService {

    @SentinelResource(
        value = "createOrder",
        blockHandler = "createOrderBlockHandler",
        fallback = "createOrderFallback"
    )
    @Override
    public CreateOrderResult createOrder(OrderRequest request) {
        // 正常下单逻辑...
    }

    // 限流/熔断时的处理
    public CreateOrderResult createOrderBlockHandler(
            OrderRequest request, BlockException ex) {
        log.warn("订单创建被限流/熔断, userId={}", request.getUserId());
        throw new BusinessException("系统繁忙,请稍后重试");
    }

    // 业务异常时的降级处理
    public CreateOrderResult createOrderFallback(
            OrderRequest request, Throwable ex) {
        log.error("订单创建降级, userId={}", request.getUserId(), ex);
        throw new BusinessException("服务暂时不可用");
    }
}

整套服务治理上线后,我再也没被半夜叫起来处理雪崩问题了。有一次数据库主库宕机切换期间,Sentinel自动把库存服务熔断了,用户的下单请求快速失败返回”请稍后重试”,而不是卡30秒超时。切换完成后流量自动恢复,整个过程用户端感知非常小。

踩坑总结与架构演进方向

最后我把踩过的几个关键坑和应对策略列一下:

坑一:事务消息的回查接口必须幂等。RocketMQ可能对同一条消息回查多次,如果你在回查方法里做了写操作(比如更新状态),就会出问题。我的做法是回查方法只读不写。

坑二:消费者重试间隔要合理配置。RocketMQ默认的重试间隔是10秒、30秒、1分钟……直到16次。库存扣减如果是因为短暂网络问题失败,10秒后重试通常能成功。但如果是因为库存真的不够,重试16次也没意义,反而会增加系统负担。我们最终把并发消费线程数设为20,重试最多5次,超过就进死信队列人工处理。

坑三:Go网关和服务之间的服务发现要统一。我们一开始Go网关和Spring Boot服务各自注册到不同的Nacos namespace,导致网关找不到服务。后来统一到同一个namespace和group后解决。

坑四:Sentinel的集群限流需要独立的Token Server。单机限流在多实例部署时效果不好,集群限流才精确。但Token Server本身也需要高可用,我们部署了3个节点做主备。

目前这套架构支撑了日均50万订单的处理量,核心链路可用性从99.5%提升到了99.95%。下一步我打算在两个方向继续演进:一是引入Seata的AT模式处理部分对一致性要求更高的场景,作为事务消息的补充;二是把Go网关的限流策略从本地令牌桶迁移到基于Redis的分布式限流,实现更精细的多租户流量管控。

微服务架构没有银弹,分布式事务也没有完美方案。但至少在我这个电商场景下,RocketMQ事务消息+幂等消费+Go网关+Sentinel熔断这套组合拳,是经过生产验证的、可用性和复杂度之间最平衡的方案。如果你也在做类似的事情,希望我的经验能帮你少走点弯路。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-jia-gou-shi-zhan-wo-shi-ru-he-yong/

(0)
小编小编
上一篇 17小时前
下一篇 17小时前

相关推荐