业务分析
一般而言,商品秒杀大概可以拆分成以下几步:
- 用户校验 校验是否多次抢单,保证每个商品每个用户只能秒杀一次
- 下单 订单信息进入消息队列,等待消费
- 减少库存 消费订单消息,减少商品库存,增加订单记录
- 付款 十五分钟内完成支付,修改支付状态
创建表
goods_info 商品库存表
列
说明
id
主键(uuid)
goods_name
商品名称
goods_stock
商品库存
1package com.jason.seckill.order.entity; 2 3/** 4 * 商品库存 5 */ 6 7public class GoodsInfo { 8 9 private String id; 10 private String goodsName; 11 private String goodsStock; 12 13 public String getId() { 14 return id; 15 } 16 17 public void setId(String id) { 18 this.id = id; 19 } 20 21 public String getGoodsName() { 22 return goodsName; 23 } 24 25 public void setGoodsName(String goodsName) { 26 this.goodsName = goodsName; 27 } 28 29 public String getGoodsStock() { 30 return goodsStock; 31 } 32 33 public void setGoodsStock(String goodsStock) { 34 this.goodsStock = goodsStock; 35 } 36 37 @Override 38 public String toString() { 39 return "GoodsInfo{" + 40 "id='" + id + '\'' + 41 ", goodsName='" + goodsName + '\'' + 42 ", goodsStock='" + goodsStock + '\'' + 43 '}'; 44 } 45} 46
order_info 订单记录表
列
说明
id
主键(uuid)
user_id
用户id
goods_id
商品id
pay_status
支付状态(0-超时未支付 1-已支付 2-待支付)
1package com.jason.seckill.order.entity; 2 3/** 4 * 下单记录 5 */ 6public class OrderRecord { 7 8 private String id; 9 private String userId; 10 private String goodsId; 11 /** 12 * 0-超时未支付 1-已支付 2-待支付 13 */ 14 private Integer payStatus; 15 16 public String getId() { 17 return id; 18 } 19 20 public void setId(String id) { 21 this.id = id; 22 } 23 24 public String getUserId() { 25 return userId; 26 } 27 28 public void setUserId(String userId) { 29 this.userId = userId; 30 } 31 32 public String getGoodsId() { 33 return goodsId; 34 } 35 36 public void setGoodsId(String goodsId) { 37 this.goodsId = goodsId; 38 } 39 40 public Integer getPayStatus() { 41 return payStatus; 42 } 43 44 public void setPayStatus(Integer payStatus) { 45 this.payStatus = payStatus; 46 } 47 48 @Override 49 public String toString() { 50 return "OrderRecord{" + 51 "id='" + id + '\'' + 52 ", userId='" + userId + '\'' + 53 ", goodsId='" + goodsId + '\'' + 54 '}'; 55 } 56} 57
功能实现
1.用户校验
使用redis做用户校验,保证每个用户每个商品只能抢一次,上代码:
1public boolean checkSeckillUser(OrderRequest order) { 2 String key = env.getProperty("seckill.redis.key.prefix") + order.getUserId() + order.getGoodsId(); 3 return redisTemplate.opsForValue().setIfAbsent(key,"1"); 4 }
userId+orderId的组合作为key,利用redis的setnx分布式锁原理来实现。如果是限时秒杀,可以通过设置key的过期时间来实现。
2.下单
下单信息肯定是要先扔到消息队列里的,这里采用RabbitMQ来做消息队列,先来看一下消息队列的模型图:
rabbitmq的配置:
1#rabbitmq配置 2spring.rabbitmq.host=127.0.0.1 3spring.rabbitmq.port=5672 4spring.rabbitmq.username=guest 5spring.rabbitmq.password=guest 6#消费者数量 7spring.rabbitmq.listener.simple.concurrency=5 8#最大消费者数量 9spring.rabbitmq.listener.simple.max-concurrency=10 10#消费者每次从队列获取的消息数量。写多了,如果长时间得不到消费,数据就一直得不到处理 11spring.rabbitmq.listener.simple.prefetch=1 12#消费接收确认机制-手动确认 13spring.rabbitmq.listener.simple.acknowledge-mode=manual 14 15mq.env=local 16#订单处理队列 17#交换机名称 18order.mq.exchange.name=${mq.env}:order:mq:exchange 19#队列名称 20order.mq.queue.name=${mq.env}:order:mq:queue 21#routingkey 22order.mq.routing.key=${mq.env}:order:mq:routing:key
rabbitmq配置类OrderRabbitmqConfig:
1/** 2 * rabbitmq配置 3 */ 4@Configuration 5public class OrderRabbitmqConfig { 6 7 private static final Logger logger = LoggerFactory.getLogger(OrderRabbitmqConfig.class); 8 9 10 @Autowired 11 private Environment env; 12 13 /** 14 * channel链接工厂 15 */ 16 @Autowired 17 private CachingConnectionFactory connectionFactory; 18 19 /** 20 * 监听器容器配置 21 */ 22 @Autowired 23 private SimpleRabbitListenerContainerFactoryConfigurer factoryConfigurer; 24 25 /** 26 * 声明rabbittemplate 27 * @return 28 */ 29 @Bean 30 public RabbitTemplate rabbitTemplate(){ 31 //消息发送成功确认,对应application.properties中的spring.rabbitmq.publisher-confirms=true 32 connectionFactory.setPublisherConfirms(true); 33 //消息发送失败确认,对应application.properties中的spring.rabbitmq.publisher-returns=true 34 connectionFactory.setPublisherReturns(true); 35 RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); 36 //设置消息发送格式为json 37 rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter()); 38 rabbitTemplate.setMandatory(true); 39 //消息发送到exchange回调 需设置:spring.rabbitmq.publisher-confirms=true 40 rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() { 41 @Override 42 public void confirm(CorrelationData correlationData, boolean ack, String cause) { 43 logger.info("消息发送成功:correlationData({}),ack({}),cause({})",correlationData,ack,cause); 44 } 45 }); 46 //消息从exchange发送到queue失败回调 需设置:spring.rabbitmq.publisher-returns=true 47 rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() { 48 @Override 49 public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) { 50 logger.info("消息丢失:exchange({}),route({}),replyCode({}),replyText({}),message:{}",exchange,routingKey,replyCode,replyText,message); 51 } 52 }); 53 return rabbitTemplate; 54 } 55 56 //---------------------------------------订单队列------------------------------------------------------ 57 58 /** 59 * 声明订单队列的交换机 60 * @return 61 */ 62 @Bean("orderTopicExchange") 63 public TopicExchange orderTopicExchange(){ 64 //设置为持久化 不自动删除 65 return new TopicExchange(env.getProperty("order.mq.exchange.name"),true,false); 66 } 67 68 /** 69 * 声明订单队列 70 * @return 71 */ 72 @Bean("orderQueue") 73 public Queue orderQueue(){ 74 return new Queue(env.getProperty("order.mq.queue.name"),true); 75 } 76 77 /** 78 * 将队列绑定到交换机 79 * @return 80 */ 81 @Bean 82 public Binding simpleBinding(){ 83 return BindingBuilder.bind(orderQueue()).to(orderTopicExchange()).with(env.getProperty("order.mq.routing.key")); 84 } 85 86 /** 87 * 注入订单对列消费监听器 88 */ 89 @Autowired 90 private OrderListener orderListener; 91 92 /** 93 * 声明订单队列监听器配置容器 94 * @return 95 */ 96 @Bean("orderListenerContainer") 97 public SimpleMessageListenerContainer orderListenerContainer(){ 98 //创建监听器容器工厂 99 SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); 100 //将配置信息和链接信息赋给容器工厂 101 factoryConfigurer.configure(factory,connectionFactory); 102 //容器工厂创建监听器容器 103 SimpleMessageListenerContainer container = factory.createListenerContainer(); 104 //指定监听器 105 container.setMessageListener(orderListener); 106 //指定监听器监听的队列 107 container.setQueues(orderQueue()); 108 return container; 109 } 110 111}
配置类声明了订单队列,交换机,通过指定的routingkey绑定了队列与交换机。另外,rabbitTemplate用来发送消息,ListenerContainer指定监听器(消费者)监听的队列。
客户下单,生产消息,上代码:
1@Service 2public class SeckillService { 3 4 private static final Logger logger = LoggerFactory.getLogger(SeckillService.class); 5 6 @Autowired 7 private RabbitTemplate rabbitTemplate; 8 @Autowired 9 private Environment env; 10 11 /** 12 * 生产消息 13 * @param order 14 */ 15 public void seckill(OrderRequest order){ 16 //设置交换机 17 rabbitTemplate.setExchange(env.getProperty("order.mq.exchange.name")); 18 //设置routingkey 19 rabbitTemplate.setRoutingKey(env.getProperty("order.mq.routing.key")); 20 //创建消息体 21 Message msg = MessageBuilder.withBody(JSON.toJSONString(order).getBytes()).build(); 22 //发送消息 23 rabbitTemplate.convertAndSend(msg); 24 } 25}
很简单,操作rabbitTemplate,指定交换机和routingkey,发送消息到绑定的队列,等待消费处理。
3.减少库存
消费者消费订单消息,做业务处理。 看一下监听器(消费者)OrderListener:
1/** 2 * 消息监听器(消费者) 3 */ 4@Component 5public class OrderListener implements ChannelAwareMessageListener { 6 7 private static final Logger logger = LoggerFactory.getLogger(OrderListener.class); 8 9 @Autowired 10 private OrderService orderService; 11 /** 12 * 处理接收到的消息 13 * @param message 消息体 14 * @param channel 通道,确认消费用 15 * @throws Exception 16 */ 17 @Override 18 public void onMessage(Message message, Channel channel) throws Exception { 19 try{ 20 //获取交付tag 21 long tag = message.getMessageProperties().getDeliveryTag(); 22 String str = new String(message.getBody(),"utf-8"); 23 logger.info("接收到的消息:{}",str); 24 JSONObject obj = JSONObject.parseObject(str); 25 //下单,操作数据库 26 orderService.order(obj.getString("userId"),obj.getString("goodsId")); 27 //确认消费 28 channel.basicAck(tag,true); 29 }catch(Exception e){ 30 logger.error("消息监听确认机制发生异常:",e.fillInStackTrace()); 31 } 32 } 33}
业务处理 OrderService:
1@Service 2public class OrderService { 3 4 @Resource 5 private SeckillMapper seckillMapper; 6 7 /** 8 * 下单,操作数据库 9 * @param userId 10 * @param goodsId 11 */ 12 @Transactional() 13 public void order(String userId,String goodsId){ 14 //该商品库存-1(当库存>0时) 15 int count = seckillMapper.reduceGoodsStockById(goodsId); 16 //更新成功,表明抢单成功,插入下单记录,支付状态设为2-待支付 17 if(count > 0){ 18 OrderRecord orderRecord = new OrderRecord(); 19 orderRecord.setId(CommonUtils.createUUID()); 20 orderRecord.setGoodsId(goodsId); 21 orderRecord.setUserId(userId); 22 orderRecord.setPayStatus(2); 23 seckillMapper.insertOrderRecord(orderRecord); 24 } 25 } 26 27}
Dao接口和Mybatis文件就不往出贴了,这里的逻辑是,update goods_info set goods_stock = goods_stock-1 where goods_stock > 0 and id=#{goodsId},这条update相当于将查询库存和减少库存合并为一个原子操作,避免高并发问题,执行成功,插入订单记录,执行失败,则库存不够抢单失败。
4.支付
订单处理完成后,如果库存减少,也就是抢单成功,那么需要用户在十五分钟内完成支付,这块就要用到死信队列(延迟队列)来处理了,先看模型图:
DLX:dead-letter Exchange 死信交换机 DLK:dead-letter RoutingKey 死信路由 ttl:time-to-live 超时时间 死信队列中,消息到期后,会通过DLX和DLK进入到pay-queue,进行消费。这是另一组消息队列,和订单消息队列是分开的。这里注意他们的绑定关系,主交换机绑定死信队列,死信交换机绑定的是主队列(pay queue)。 接下来声明图中的一系列组件,首先application.properties中增加配置:
1#支付处理队列 2#主交换机 3pay.mq.exchange.name=${mq.env}:pay:mq:exchange 4#死信交换机(DLX) 5pay.dead-letter.mq.exchange.name=${mq.env}:pay:dead-letter:mq:exchange 6#主队列 7pay.mq.queue.name=${mq.env}:pay:mq:queue 8#死信队列 9pay.dead-letter.mq.queue.name=${mq.env}:pay:dead-letter:mq:queue 10#主routingkey 11pay.mq.routing.key=${mq.env}:pay:mq:routing:key 12#死信routingkey(DLK) 13pay.dead-letter.mq.routing.key=${mq.env}:pay:dead-letter:mq:routing:key 14#支付超时时间(毫秒)(TTL),测试原因,这里模拟5秒,如果是生产环境,这里可以是15分钟等 15pay.mq.ttl=5000
配置类OrderRabbitmqConfig中增加支付队列和死信队列的声明:
1 /** 2 * 死信队列,十五分钟超时 3 * @return 4 */ 5 @Bean 6 public Queue payDeadLetterQueue(){ 7 Map args = new HashMap(); 8 //声明死信交换机 9 args.put("x-dead-letter-exchange",env.getProperty("pay.dead-letter.mq.exchange.name")); 10 //声明死信routingkey 11 args.put("x-dead-letter-routing-key",env.getProperty("pay.dead-letter.mq.routing.key")); 12 //声明死信队列中的消息过期时间 13 args.put("x-message-ttl",env.getProperty("pay.mq.ttl",int.class)); 14 //创建死信队列 15 return new Queue(env.getProperty("pay.dead-letter.mq.queue.name"),true,false,false,args); 16 } 17 18 /** 19 * 支付队列交换机(主交换机) 20 * @return 21 */ 22 @Bean 23 public TopicExchange payTopicExchange(){ 24 return new TopicExchange(env.getProperty("pay.mq.exchange.name"),true,false); 25 } 26 27 /** 28 * 将主交换机绑定到死信队列 29 * @return 30 */ 31 @Bean 32 public Binding payBinding(){ 33 return BindingBuilder.bind(payDeadLetterQueue()).to(payTopicExchange()).with(env.getProperty("pay.mq.routing.key")); 34 } 35 36 /** 37 * 支付队列(主队列) 38 * @return 39 */ 40 @Bean 41 public Queue payQueue(){ 42 return new Queue(env.getProperty("pay.mq.queue.name"),true); 43 } 44 45 /** 46 * 死信交换机 47 * @return 48 */ 49 @Bean 50 public TopicExchange payDeadLetterExchange(){ 51 return new TopicExchange(env.getProperty("pay.dead-letter.mq.exchange.name"),true,false); 52 } 53 54 /** 55 * 将主队列绑定到死信交换机 56 * @return 57 */ 58 @Bean 59 public Binding payDeadLetterBinding(){ 60 return BindingBuilder.bind(payQueue()).to(payDeadLetterExchange()).with(env.getProperty("pay.dead-letter.mq.routing.key")); 61 } 62 63 /** 64 * 注入支付监听器 65 */ 66 @Autowired 67 private PayListener payListener; 68 69 /** 70 * 支付队列监听器容器 71 * @return 72 */ 73 @Bean 74 public SimpleMessageListenerContainer payMessageListenerContainer(){ 75 SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); 76 factoryConfigurer.configure(factory,connectionFactory); 77 SimpleMessageListenerContainer listenerContainer = factory.createListenerContainer(); 78 listenerContainer.setMessageListener(payListener); 79 listenerContainer.setQueues(payQueue()); 80 return listenerContainer; 81 }
支付队列和死信队列的Queue、Exchange、routingkey都已就绪。 看生产者:
1@Service 2public class OrderService { 3 4 @Resource 5 private SeckillMapper seckillMapper; 6 7 @Autowired 8 private RabbitTemplate rabbitTemplate; 9 10 @Autowired 11 private Environment env; 12 13 /** 14 * 下单,操作数据库 15 * @param userId 16 * @param goodsId 17 */ 18 @Transactional() 19 public void order(String userId,String goodsId){ 20 //该商品库存-1(当库存>0时) 21 int count = seckillMapper.reduceGoodsStockById(goodsId); 22 //更新成功,表明抢单成功,插入下单记录,支付状态设为2-待支付 23 if(count > 0){ 24 OrderRecord orderRecord = new OrderRecord(); 25 orderRecord.setId(CommonUtils.createUUID()); 26 orderRecord.setGoodsId(goodsId); 27 orderRecord.setUserId(userId); 28 orderRecord.setPayStatus(2); 29 seckillMapper.insertOrderRecord(orderRecord); 30 //将该订单添加到支付队列 31 rabbitTemplate.setExchange(env.getProperty("pay.mq.exchange.name")); 32 rabbitTemplate.setRoutingKey(env.getProperty("pay.mq.routing.key")); 33 rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter()); 34 String json = JSON.toJSONString(orderRecord); 35 Message msg = MessageBuilder.withBody(json.getBytes()).build(); 36 rabbitTemplate.convertAndSend(msg); 37 } 38 } 39}
在OrderService中,数据库操作完成后,将订单信息发送到死信队列,死信队列中的消息会在十五分钟后进入到支付队列,等待消费。 再看消费者:
1@Component 2public class PayListener implements ChannelAwareMessageListener { 3 4 private static final Logger logger = LoggerFactory.getLogger(PayListener.class); 5 6 @Autowired 7 private PayService payService; 8 9 @Override 10 public void onMessage(Message message, Channel channel) throws Exception { 11 Long tag = message.getMessageProperties().getDeliveryTag(); 12 try { 13 String str = new String(message.getBody(), "utf-8"); 14 logger.info("接收到的消息:{}",str); 15 JSONObject json = JSON.parseObject(str); 16 String orderId = json.getString("id"); 17 //确认是否付款 18 payService.confirmPay(orderId); 19 //确认消费 20 channel.basicAck(tag, true); 21 }catch(Exception e){ 22 logger.info("支付消息消费出错:{}",e.getMessage()); 23 logger.info("出错的tag:{}",tag); 24 } 25 } 26}
PayService:
1@Service 2public class PayService { 3 4 private static final Logger logger = LoggerFactory.getLogger(PayService.class); 5 6 @Resource 7 private SeckillMapper seckillMapper; 8 9 /** 10 * 确认是否支付 11 * @param orderId 12 */ 13 public void confirmPay(String orderId){ 14 OrderRecord orderRecord = seckillMapper.selectNoPayOrderById(orderId); 15 //根据订单号校验该用户是否已支付 16 if(checkPay(orderId)){ 17 //已支付 18 orderRecord.setPayStatus(1); 19 seckillMapper.updatePayStatus(orderRecord); 20 logger.info("用户{}已支付",orderId); 21 }else{ 22 //未支付 23 orderRecord.setPayStatus(0); 24 seckillMapper.updatePayStatus(orderRecord); 25 //取消支付后,商品库存+1 26 seckillMapper.returnStock(orderRecord.getGoodsId()); 27 logger.info("用户{}未支付",orderId); 28 } 29 } 30 31 /** 32 * 模拟判断订单支付成功或失败,成功失败随机 33 * @param orderId 34 * @return 35 */ 36 public boolean checkPay(String orderId){ 37 Random random = new Random(); 38 int res = random.nextInt(2); 39 return res==0?false:true; 40 }
这里checkPay()方法模拟调用第三方支付接口来判断用户是否已支付。若支付成功,订单改为已支付状态,支付失败,改为已取消状态,库存退回。
总结
整个demo,是两组消息队列撑起来的,一组订单消息队列,一组支付消息队列,而每一组队列都是由queue、exchange、routingkey、生产者以及消费者组成。交换机通过routingkey绑定队列,rabbitTemplate通过指定交换机和routingkey将消息发送到指定队列,消费者监听该队列进行消费。不同的是第二组支付队列里嵌入了死信队列来做一个十五分钟的延迟支付。