电商返利系统高并发解决方案:Java 分布式锁、限流与削峰实践
电商返利系统高并发解决方案: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 研发团队,转载请注明出处!
更多推荐



所有评论(0)