# stage06-03-rocketmq **Repository Path**: null_631_9084/stage06-03-rocketmq ## Basic Information - **Project Name**: stage06-03-rocketmq - **Description**: No description available - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2020-11-20 - **Last Updated**: 2020-12-19 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README 基于RocketMQ设计秒杀。 要求: 1. 秒杀商品LagouPhone,数量100个。 2. 秒杀商品不能超卖。 3. 抢购链接隐藏 4. Nginx+Redis+RocketMQ+Tomcat+MySQL ### 配置 ```properties spring.application.name=rocket-mq-order server.port=9999 spring.redis.host=localhost spring.redis.database=0 spring.redis.port=6379 rocketmq.name-server=192.168.181.141:9876 rocketmq.producer.group=producer_order_gr01 ``` * RedisConfig ```java package com.liu.rocketmq.redis; import com.alibaba.fastjson.support.spring.GenericFastJsonRedisSerializer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.CacheManager; import org.springframework.cache.annotation.EnableCaching; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.redis.cache.RedisCacheConfiguration; import org.springframework.data.redis.cache.RedisCacheManager; import org.springframework.data.redis.connection.jedis.JedisConnectionFactory; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.serializer.RedisSerializationContext; import org.springframework.data.redis.serializer.StringRedisSerializer; import java.time.Duration; @EnableCaching @Configuration public class RedisConfig { @Autowired private JedisConnectionFactory redisConnectionFactory; @Bean public CacheManager cacheManager() { // StringRedisSerializer keySerializer = new StringRedisSerializer(); GenericFastJsonRedisSerializer valueSerializer = new GenericFastJsonRedisSerializer(); RedisCacheConfiguration config = RedisCacheConfiguration.defaultCacheConfig(); config = config.serializeValuesWith(RedisSerializationContext.SerializationPair.fromSerializer(valueSerializer)) .entryTtl(Duration.ofHours(2)) .prefixKeysWith("order.rocketmq"); RedisCacheManager cacheManager = RedisCacheManager.builder(redisConnectionFactory) .cacheDefaults(config) .build(); return cacheManager; } @Bean public RedisTemplate redisTemplate() { // 配置redisTemplate RedisTemplate redisTemplate = new RedisTemplate<>(); redisTemplate.setConnectionFactory(redisConnectionFactory); StringRedisSerializer stringSerializer = new StringRedisSerializer(); GenericFastJsonRedisSerializer valueSerializer = new GenericFastJsonRedisSerializer(); redisTemplate.setKeySerializer(stringSerializer); // key序列化 redisTemplate.setHashKeySerializer(stringSerializer); redisTemplate.setValueSerializer(valueSerializer); redisTemplate.setHashValueSerializer(valueSerializer); redisTemplate.afterPropertiesSet(); return redisTemplate; } } ``` * RedissonManager-redisson对象 ```java public class RedissonManager { private static Config config = new Config(); //声明redisso对象 private static RedissonClient redisson = null; //实例化redisson static { config.useSingleServer() .setAddress("redis://localhost:6379") .setDatabase(0); //得到redisson对象 //redisson =(Redisson) Redisson.create(config); redisson= Redisson.create(config); } //获取redisson对象的方法 public static RedissonClient getRedisson(){ return redisson; } } ``` * 分布式锁 ```java package com.liu.rocketmq.redis; import org.redisson.Redisson; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import java.util.concurrent.TimeUnit; public class DistributedRedisLock { //从配置类中获取redisson对象 private static RedissonClient redisson = RedissonManager.getRedisson(); private static final String LOCK_TITLE = "redisLock_"; //加锁 public static boolean lock(String lockName, Integer lockTime, TimeUnit timeUnit) { String key = LOCK_TITLE + lockName; RLock mylock = redisson.getLock(key); mylock.lock(lockTime, timeUnit); return true; } //锁的释放 public static void unLock(String lockName) { String key = LOCK_TITLE + lockName; RLock mylock = redisson.getLock(key); mylock.unlock(); } } ``` * 调用接口 对商品加锁 -把 商品-数量100 设置到redis中 ```java @Slf4j @RestController public class OrderController { private static final String goodsName = "LagouPhone"; @Autowired private RedisTemplate redisTemplate; @Autowired private RocketMQTemplate rocketMQTemplate; @GetMapping("/setGoodsInvent") public ResponseEntity setGoodsInventory() { /** * 加锁-商品库存放入Redis中 */ DistributedRedisLock.lock(goodsName, 2, TimeUnit.SECONDS); //100个库存 redisTemplate.opsForValue().set(goodsName, 100); DistributedRedisLock.unLock(goodsName); return ResponseEntity.ok().build(); } } ``` * 实体 ```java public class Order implements Serializable { private String orderNo; private String goods; private Integer goodNum; } ``` * redisson分布式锁-生成订单 ```java @Slf4j @RestController public class OrderController { private static final String goodsName = "LagouPhone"; @Autowired private RedisTemplate redisTemplate; @Autowired private RocketMQTemplate rocketMQTemplate; @RequestMapping("/order") public ResponseEntity order() { ExecutorService executorService = Executors.newCachedThreadPool(); for (int i = 0; i < 1000; i++) { executorService.execute(() -> { if (DistributedRedisLock.lock(goodsName, 2, TimeUnit.SECONDS)) { if (redisTemplate.opsForValue().get(goodsName) > 0) { //生成记录 秒杀就只有一个 Order order = new Order(); order.setGoods(goodsName); order.setGoodNum(1); order.setOrderNo(UUID.randomUUID().toString()); //库存减一 redisTemplate.opsForValue().decrement(goodsName,1); rocketMQTemplate.convertAndSend("order", order); } else { log.error("库存不足"); } DistributedRedisLock.unLock(goodsName); } }); } return ResponseEntity.ok().build(); } } ``` * 订单消费 ```java @Slf4j @Component @RocketMQMessageListener(topic = "order", consumerGroup = "order_group") public class OrderConsumer implements RocketMQListener { @Autowired private RocketMQTemplate rocketMQTemplate; @Override public void onMessage(Order order) { //模拟订单业务处理 log.info("接受到下单请求:"+order); log.info("下单业务处理"); log.info("等待用户支付"); rocketMQTemplate.syncSend("pay_order", MessageBuilder.withPayload(order).build(), 2000, 1); } } ``` * 支付消费者 ```java @Slf4j @Component @RocketMQMessageListener(topic = "pay_order", consumerGroup = "pay_group") public class PayConsumer implements RocketMQListener { @Autowired private RocketMQTemplate rocketMQTemplate; @Override public void onMessage(Order order) { //模拟订单业务处理 log.info("订单支付状态检查:"+order.toString()); log.info("用户支付成功"); } } ```