电商返利系统高并发解决方案:Java 分布式锁、限流与削峰实践

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

在“双11”、“618”等电商大促期间,返利系统瞬间面临的流量洪峰是平日的数十倍甚至上百倍。用户集中点击下单、联盟接口回调激增、佣金计算请求堆积,任何环节的处理不当都可能导致数据错乱(如重复发佣)、服务雪崩或资金损失。对于省赚客APP而言,保障高并发下的数据一致性与系统可用性是核心命题。本文将深入探讨如何利用Java生态中的分布式锁、多维限流算法以及消息队列削峰填谷机制,构建一套坚不可摧的高并发防御体系。

基于Redisson的分布式锁与幂等性保障

在分布式环境下,同一用户的订单可能被多次回调,或者用户快速连续点击领取奖励,极易引发并发竞争条件(Race Condition)。传统的synchronized仅能锁住单机进程,无法解决集群环境下的资源争抢。我们采用基于Redis的Redisson分布式锁,结合业务唯一键(如订单号+用户ID),确保核心逻辑(如佣金入账、库存扣减)在同一时刻只有一个线程执行,从而实现严格的幂等性。

package juwatech.cn.concurrency.lock;

import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import juwatech.cn.service.CommissionService;
import juwatech.cn.model.OrderContext;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.TimeUnit;

@Slf4j
public class DistributedLockExecutor {

    private final RedissonClient redissonClient;
    private final CommissionService commissionService;

    public DistributedLockExecutor(RedissonClient redissonClient, CommissionService commissionService) {
        this.redissonClient = redissonClient;
        this.commissionService = commissionService;
    }

    /**
     * 执行带分布式锁的佣金结算逻辑
     */
    public void executeWithLock(OrderContext context) {
        String lockKey = "lock:commission:" + context.getUserId() + ":" + context.getOrderId();
        RLock lock = redissonClient.getLock(lockKey);

        boolean isLocked = false;
        try {
            // 尝试获取锁,等待5秒,自动释放时间10秒(看门狗机制会自动续期)
            isLocked = lock.tryLock(5, 10, TimeUnit.SECONDS);
            
            if (isLocked) {
                // 1. 二次检查幂等性(防止锁释放瞬间的重复请求)
                if (commissionService.isAlreadyProcessed(context.getOrderId())) {
                    log.info("Order already processed, skipping: {}", context.getOrderId());
                    return;
                }
                
                // 2. 执行核心业务逻辑
                commissionService.calculateAndCredit(context);
                log.info("Commission settled successfully for order: {}", context.getOrderId());
            } else {
                log.warn("Failed to acquire lock, request rejected: {}", context.getOrderId());
                // 可选择不处理或放入重试队列
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            log.error("Lock acquisition interrupted", e);
        } finally {
            if (isLocked && lock.isHeldByCurrentThread()) {
                lock.unlock();
            }
        }
    }
}

多层级限流策略:网关拦截与服务保护

面对海量请求,第一道防线必须是限流。我们在API网关层(基于Spring Cloud Gateway)实施全局限流,防止恶意刷接口;在服务层针对核心方法实施细粒度限流,保护数据库和下游联盟接口。利用Sentinel实现基于QPS、线程数及热点参数的多维控制,一旦触发阈值,立即执行降级逻辑,返回友好提示或默认值,避免系统过载崩溃。

package juwatech.cn.concurrency.limit;

import com.alibaba.csp.sentinel.annotation.SentinelResource;
import com.alibaba.csp.sentinel.slots.block.BlockException;
import com.alibaba.csp.sentinel.slots.block.flow.FlowRule;
import com.alibaba.csp.sentinel.slots.block.flow.FlowRuleManager;
import juwatech.cn.dto.UserRewardRequest;
import juwatech.cn.dto.UserRewardResponse;
import lombok.extern.slf4j.Slf4j;
import javax.annotation.PostConstruct;
import java.util.Collections;

@Slf4j
public class RateLimitService {

    @PostConstruct
    public void initRules() {
        // 配置规则:针对"claimReward"资源,QPS阈值为1000
        FlowRule rule = new FlowRule("claimReward")
                .setCount(1000)
                .setGrade(1) // QPS模式
                .setLimitApp("default");
        FlowRuleManager.loadRules(Collections.singletonList(rule));
        log.info("Rate limiting rules initialized");
    }

    /**
     * 领取奖励接口,受Sentinel保护
     */
    @SentinelResource(value = "claimReward", blockHandler = "handleBlock")
    public UserRewardResponse claimReward(UserRewardRequest request) {
        // 模拟耗时业务操作
        juwatech.cn.service.RewardProcessor.process(request);
        return new UserRewardResponse(true, "Success");
    }

    /**
     * 限流后的兜底处理方法
     */
    public UserRewardResponse handleBlock(UserRewardRequest request, BlockException ex) {
        log.warn("Request blocked due to high concurrency: {}", request.getUserId());
        // 返回降级响应,引导用户稍后重试
        return new UserRewardResponse(false, "System busy, please try again later.");
    }
}

消息队列削峰填谷与异步解耦

大促期间的流量具有极强的突发性,直接同步处理会导致数据库连接池瞬间耗尽。我们引入RocketMQ作为缓冲层,将非实时的核心业务(如订单同步、佣金计算、积分发放)异步化。用户请求到达后,仅需将消息写入队列即可快速响应,后端消费者根据自身处理能力匀速拉取消息进行消费。这种“削峰填谷”机制将瞬时洪峰转化为平稳的流量流,极大提升了系统的吞吐量。

package juwatech.cn.concurrency.async;

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import juwatech.cn.model.OrderSyncMessage;
import juwatech.cn.repository.OrderBufferRepository;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;

@Component
@RequiredArgsConstructor
public class OrderAsyncProcessor {

    private final RocketMQTemplate rocketMQTemplate;
    private final OrderBufferRepository bufferRepository;

    private static final String TOPIC_ORDER_SYNC = "TOPIC_ORDER_SYNC_HIGH_CONCURRENCY";

    /**
     * 接收上游回调,快速写入消息队列
     */
    public void receiveUpstreamCallback(String orderId, String platformData) {
        // 1. 极简校验后直接发送消息,不执行任何DB写操作
        OrderSyncMessage message = new OrderSyncMessage(orderId, platformData, System.currentTimeMillis());
        
        // 发送顺序消息或普通消息,确保高吞吐
        rocketMQTemplate.sendOneWay(TOPIC_ORDER_SYNC, message);
        
        // 2. 可选:记录原始日志到高速存储(如HBase/MongoDB)用于审计
        bufferRepository.saveRawLog(orderId, platformData);
    }
}

// 消费者端代码
package juwatech.cn.consumer;

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import juwatech.cn.model.OrderSyncMessage;
import juwatech.cn.service.CommissionCalcService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;

@Slf4j
@Service
@RocketMQMessageListener(topic = "TOPIC_ORDER_SYNC_HIGH_CONCURRENCY", consumerGroup = "CG_ORDER_CONSUMER")
public class OrderSyncConsumer implements RocketMQListener<OrderSyncMessage> {

    private final CommissionCalcService calcService;

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

    @Override
    public void onMessage(OrderSyncMessage message) {
        try {
            // 3. 消费者按照自身能力匀速处理
            // 即使上游每秒进来1万单,消费者也可以控制在每秒处理2000单
            calcService.processOrder(message.getOrderId(), message.getPlatformData());
            log.info("Order processed asynchronously: {}", message.getOrderId());
        } catch (Exception e) {
            log.error("Failed to process order, will retry via MQ mechanism", e);
            throw new RuntimeException(e); // 触发MQ重试
        }
    }
}

数据库连接池优化与读写分离

在高并发场景下,数据库往往是最后的瓶颈。除了应用层的优化,我们还对数据库架构进行了深度改造。实施主从读写分离,将大量的查询请求(如用户账单查询、商品详情)路由到从库,主库专注于写入。同时,精细化配置HikariCP连接池参数,根据CPU核数和业务类型调整最大连接数、最小空闲连接数及超时时间,避免连接泄露和等待过长。

package juwatech.cn.config.datasource;

import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import javax.sql.DataSource;

@Configuration
public class DatabasePoolConfig {

    @Bean(name = "writerDataSource")
    public DataSource writerDataSource() {
        HikariConfig config = new HikariConfig();
        config.setJdbcUrl("jdbc:mysql://master-db:3306/shengzhuanke?useSSL=false");
        config.setUsername("root");
        config.setPassword("secret");
        
        // 高并发写场景优化
        config.setMaximumPoolSize(50); // 根据压测结果设定
        config.setMinimumIdle(10);
        config.setConnectionTimeout(30000);
        config.setIdleTimeout(600000);
        config.setMaxLifetime(1800000);
        config.addDataSourceProperty("cachePrepStmts", "true");
        config.addDataSourceProperty("prepStmtCacheSize", "250");
        config.addDataSourceProperty("prepStmtCacheSqlLimit", "2048");
        
        return new HikariDataSource(config);
    }
    
    // 读数据源配置类似,最大连接数可适当调大
}

通过分布式锁保障数据强一致性,多层限流构筑系统防火墙,消息队列实现流量平滑,以及数据库层面的深度优化,省赚客APP成功经受住了多次亿级流量大促的考验。这套高并发解决方案不仅保障了业务的连续稳定,更为用户提供了丝滑流畅的返利体验,成为平台核心竞争力的重要组成部分。

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

Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐