天猫返利 APP 高可用架构:分布式事务、最终一致性与数据可靠性

大家好,我是高佣返利省赚客 APP 研发者阿宝!

在电商返利领域,天猫(Tmall)以其庞大的 GMV 和复杂的营销规则占据核心地位。然而,对接天猫联盟意味着要面对极高的并发流量、严苛的数据一致性要求以及跨服务调用的不确定性。用户下单后,订单状态需从“创建”流转到“付款”、“结算”,期间涉及佣金预估、多级分润、钱包入账等多个微服务。任何环节的网络抖动或逻辑异常都可能导致“有订单无佣金”或“重复发钱”的严重事故。为此,省赚客 APP 构建了一套基于最终一致性原则的高可用架构,通过分布式事务控制、可靠消息驱动及多重数据校验机制,确保每一分返利都精准无误。

基于 RocketMQ 事务消息的最终一致性方案

在微服务架构下,本地事务无法跨越服务边界。传统的 TCC 或 Seata AT 模式在高并发场景下锁资源开销大,易造成性能瓶颈。我们采用 RocketMQ 事务消息机制,将“订单落库”与“发送佣金计算消息”这两个操作绑定为原子单元。只有当本地事务成功提交,消息才会发送给下游;若本地回滚,消息自动丢弃。下游服务消费消息实现异步解耦,即使短暂不可用,消息也会持久化等待重试,从而保证数据的最终一致性。

package juwatech.cn.tmall.transaction;

import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import juwatech.cn.order.entity.TmallOrder;
import juwatech.cn.order.repository.TmallOrderRepository;
import lombok.extern.slf4j.Slf4j;
import org.springframework.messaging.Message;
import java.util.Optional;

@Slf4j
@RocketMQTransactionListener(txProducerGroup = "pg_tmall_order_tx")
public class TmallOrderTransactionListener implements RocketMQLocalTransactionListener {

    private final TmallOrderRepository orderRepository;

    public TmallOrderTransactionListener(TmallOrderRepository orderRepository) {
        this.orderRepository = orderRepository;
    }

    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        // 1. 解析消息中的订单号
        String tradeId = parseTradeId(msg);
        
        // 2. 检查本地数据库订单状态
        Optional<TmallOrder> orderOpt = orderRepository.findByTradeId(tradeId);
        
        if (orderOpt.isPresent()) {
            TmallOrder order = orderOpt.get();
            if ("CREATED".equals(order.getStatus())) {
                // 本地事务成功,确认提交消息,触发下游佣金计算
                log.info("Local transaction success, commit message for order: {}", tradeId);
                return RocketMQLocalTransactionState.COMMIT;
            } else if ("INVALID".equals(order.getStatus())) {
                // 订单已失效,回滚消息
                return RocketMQLocalTransactionState.ROLLBACK;
            }
        }
        
        // 状态未知,返回 UNKNOWN,等待 MQ 后台回查
        log.warn("Local transaction status unknown, return UNKNOWN for order: {}", tradeId);
        return RocketMQLocalTransactionState.UNKNOWN;
    }

    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
        // 3. 事务回查逻辑:当 Broker 未收到 Confirm/Rollback 时触发
        String tradeId = parseTradeId(msg);
        Optional<TmallOrder> orderOpt = orderRepository.findByTradeId(tradeId);
        
        if (orderOpt.isPresent() && "CREATED".equals(orderOpt.get().getStatus())) {
            log.info("Transaction check passed, committing message for: {}", tradeId);
            return RocketMQLocalTransactionState.COMMIT;
        }
        
        log.info("Transaction check failed, rolling back message for: {}", tradeId);
        return RocketMQLocalTransactionState.ROLLBACK;
    }

    private String parseTradeId(Message msg) {
        // 解析 Payload 逻辑
        return "TMALL_123456"; 
    }
}

分布式幂等性设计与防重放机制

天猫联盟的回调接口在网络不稳定时可能重复推送同一笔订单。为了防止用户钱包被多次注入佣金,我们在消费者端构建了多层幂等防线。第一层利用 Redis 原子操作进行快速拦截,第二层依赖数据库唯一索引强约束,第三层通过业务状态机判断,确保同一订单在同一状态下仅被执行一次。

package juwatech.cn.tmall.idempotent;

import org.springframework.data.redis.core.StringRedisTemplate;
import juwatech.cn.commission.repository.CommissionRecordRepository;
import juwatech.cn.exception.DuplicateProcessException;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;
import java.util.concurrent.TimeUnit;

@Component
@RequiredArgsConstructor
public class IdempotentProcessor {

    private final StringRedisTemplate redisTemplate;
    private final CommissionRecordRepository recordRepository;
    private static final String KEY_PREFIX = "tmall:process:";

    /**
     * 执行幂等检查并处理业务
     */
    public void processWithIdempotency(String tradeId, String status, Runnable businessLogic) {
        String lockKey = KEY_PREFIX + tradeId + ":" + status;
        
        // 1. Redis 预检:SETNX 尝试加锁
        Boolean isAbsent = redisTemplate.opsForValue().setIfAbsent(lockKey, "1", 24, TimeUnit.HOURS);
        if (Boolean.FALSE.equals(isAbsent)) {
            log.warn("Duplicate request detected by Redis for order: {}, status: {}", tradeId, status);
            return; // 直接忽略,视为已成功处理
        }

        try {
            // 2. 数据库唯一索引二次校验
            if (recordRepository.existsByTradeIdAndStatus(tradeId, status)) {
                log.info("Duplicate request detected by DB for order: {}", tradeId);
                return;
            }

            // 3. 执行业务逻辑(佣金计算、入账)
            businessLogic.run();
            
            // 4. 记录处理日志(利用唯一索引防止并发写入)
            recordRepository.saveLog(tradeId, status);
            
        } catch (Exception e) {
            log.error("Business logic failed, will retry via MQ", e);
            throw e; // 抛出异常触发 MQ 重试
        }
    }
}

全链路数据对账与自动补偿自愈

即便有完善的事务和幂等机制,极端情况(如代码 Bug、第三方数据修正)仍可能导致数据不一致。我们建立了 T+1 日的自动化对账系统,每日凌晨拉取天猫联盟官方结算报表,与本地数据库进行全量比对。针对发现的差异(长款、短款、金额不符),系统根据预设策略自动执行补偿任务,无需人工干预。

package juwatech.cn.tmall.reconciliation;

import juwatech.cn.client.TmallBillClient;
import juwatech.cn.repository.LocalCommissionRepository;
import juwatech.cn.entity.ReconciliationDiff;
import juwatech.cn.service.AutoCompensateService;
import lombok.extern.slf4j.Slf4j;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;

@Slf4j
public class DailyReconciliationTask {

    private final TmallBillClient billClient;
    private final LocalCommissionRepository localRepo;
    private final AutoCompensateService compensateService;

    public void execute(LocalDate date) {
        log.info("Starting Tmall reconciliation for date: {}", date);

        // 1. 下载官方账单
        List<TmallBillItem> remoteBills = billClient.downloadSettlementReport(date);
        
        // 2. 查询本地记录
        List<LocalCommissionRecord> localRecords = localRepo.findBySettleDate(date);
        Map<String, LocalCommissionRecord> localMap = localRecords.stream()
            .collect(Collectors.toMap(LocalCommissionRecord::getTradeId, r -> r));

        for (TmallBillItem remote : remoteBills) {
            LocalCommissionRecord local = localMap.remove(remote.getTradeId());
            
            if (local == null) {
                // 场景:本地漏单
                handleMissingOrder(remote);
            } else {
                // 场景:金额不一致(允许 0.01 误差)
                BigDecimal diff = remote.getCommission().subtract(local.getCommission());
                if (diff.abs().compareTo(new BigDecimal("0.01")) > 0) {
                    ReconciliationDiff diffRecord = new ReconciliationDiff(
                        remote.getTradeId(), 
                        "AMOUNT_MISMATCH", 
                        local.getCommission(), 
                        remote.getCommission()
                    );
                    compensateService.adjustBalance(diffRecord);
                }
                
                // 场景:状态不一致(如本地未失效,平台已退款)
                if (!remote.isValid() && local.isValid()) {
                    compensateService.reverseCommission(local.getId());
                }
            }
        }
        
        // 处理本地多出的订单(可能是虚假订单或延迟未到)
        for (LocalCommissionRecord orphan : localMap.values()) {
            log.warn("Local order exists but not in Tmall bill: {}", orphan.getTradeId());
            // 标记待核查,不立即冲正,防止误杀
            orphan.setVerifyStatus("PENDING_CHECK");
            localRepo.update(orphan);
        }
    }

    private void handleMissingOrder(TmallBillItem remote) {
        // 自动补单:反向生成订单记录并触发分润
        compensateService.backfillOrder(remote);
    }
}

高可用容灾与降级策略

面对天猫大促期间的流量洪峰,系统必须具备自我保护和降级能力。我们利用 Sentinel 配置了多维度的限流熔断规则。当检测到天猫 API 响应超时率飙升或下游佣金服务负载过高时,自动触发降级策略:暂停非核心的实时预估功能,将订单请求暂存至本地队列,待系统恢复后再异步处理,确保核心交易链路不崩塌。

package juwatech.cn.tmall.resilience;

import com.alibaba.csp.sentinel.annotation.SentinelResource;
import com.alibaba.csp.sentinel.slots.block.BlockException;
import juwatech.cn.service.CommissionCalcService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;

@Slf4j
@Service
public class TmallFallbackService {

    private final CommissionCalcService calcService;

    public TmallFallbackService(CommissionCalcService calcService) {
        this.calcService = calcService;
    }

    @SentinelResource(value = "calcTmallCommission", blockHandler = "handleBlock", fallback = "handleFallback")
    public void calculateCommission(String tradeId, Double amount) {
        calcService.execute(tradeId, amount);
    }

    // 限流或熔断时的兜底逻辑
    public void handleBlock(String tradeId, Double amount, BlockException ex) {
        log.warn("Triggered limit/fuse for order: {}, reason: {}", tradeId, ex.getRule().getResource());
        // 策略:写入延迟队列,稍后重试
        juwatech.cn.mq.DelayQueueProducer.sendRetryMessage(tradeId, amount, 60);
    }

    // 业务异常时的兜底逻辑
    public void handleFallback(String tradeId, Double amount, Throwable t) {
        log.error("Business exception for order: {}", tradeId, t);
        // 策略:记录错误日志,转入人工审核队列
        juwatech.cn.service.ErrorLogService.save(tradeId, t.getMessage());
    }
}

通过分布式事务消息保障数据最终一致性,多层幂等机制杜绝重复计算,自动化对账系统实现自我愈合,以及完善的容灾降级策略,省赚客 APP 构建了坚不可摧的天猫返利高可用架构。这不仅保障了海量交易下的资金安全,更为用户提供了稳定、可信的返利体验,确立了我们在行业内的技术领先地位。

本文著作权归 省赚客 app 研发团队,转载请注明出处!

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐