Jelajahi Sumber

切换为rabbitmq,目前已封装完成

mac 2 tahun lalu
induk
melakukan
b1d2e7b987
19 mengubah file dengan 819 tambahan dan 206 penghapusan
  1. 5 0
      kxs-common/kxs-common-core/src/main/java/com/kxs/common/core/constant/CommonConstants.java
  2. 7 0
      kxs-common/kxs-common-core/src/main/java/com/kxs/common/core/util/SpringContextHolder.java
  3. 6 2
      kxs-common/kxs-common-mq/pom.xml
  4. 140 0
      kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/config/RabbitConfiguration.java
  5. 79 0
      kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/config/RabbitMQBeanProcessor.java
  6. 18 0
      kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/enums/MQSendTypeEnum.java
  7. 21 0
      kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/handlers/IMQSender.java
  8. 49 0
      kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/handlers/rabbit/RabbitMQSender.java
  9. 23 0
      kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/model/AbstractMQ.java
  10. 3 1
      kxs-common/kxs-common-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports
  11. 4 0
      kxs-product/kxs-product-biz/pom.xml
  12. 4 0
      kxs-system/kxs-system-biz/pom.xml
  13. 69 0
      kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/amqp/RabbitOrderQueueListener.java
  14. 111 0
      kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/amqp/rabbit/RabbitOrderQueueMQ.java
  15. 4 0
      kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/mapper/KxsCampMapper.java
  16. 2 1
      kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/service/KxsCampService.java
  17. 245 201
      kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/service/impl/KxsCampServiceImpl.java
  18. 1 1
      kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/task/KxsSystemTask.java
  19. 28 0
      kxs-system/kxs-system-biz/src/main/resources/mapper/KxsCampMapper.xml

+ 5 - 0
kxs-common/kxs-common-core/src/main/java/com/kxs/common/core/constant/CommonConstants.java

@@ -54,6 +54,11 @@ public interface CommonConstants {
 	 */
 	String BACK_END_PROJECT = "kxs";
 
+	/**
+	 * 包名
+	 */
+	String PACKAGE_NAME = "com.kxs";
+
 	/**
 	 * 成功标记
 	 */

+ 7 - 0
kxs-common/kxs-common-core/src/main/java/com/kxs/common/core/util/SpringContextHolder.java

@@ -52,6 +52,13 @@ public class SpringContextHolder implements ApplicationContextAware, DisposableB
 		return applicationContext.getBean(requiredType);
 	}
 
+	/**
+	 * 通过name,以及Clazz返回指定的Bean
+	 */
+	public static <T> T getBean(String name, Class<T> clazz){
+		return applicationContext.getBean(name, clazz);
+	}
+
 	/**
 	 * 清除SpringContextHolder中的ApplicationContext为Null.
 	 */

+ 6 - 2
kxs-common/kxs-common-mq/pom.xml

@@ -19,8 +19,12 @@
     </properties>
     <dependencies>
         <dependency>
-            <groupId>org.apache.rocketmq</groupId>
-            <artifactId>rocketmq-spring-boot-starter</artifactId>
+            <groupId>org.springframework.boot</groupId>
+            <artifactId>spring-boot-starter-amqp</artifactId>
+        </dependency>
+        <dependency>
+            <groupId>com.kxs</groupId>
+            <artifactId>kxs-common-core</artifactId>
         </dependency>
     </dependencies>
 </project>

+ 140 - 0
kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/config/RabbitConfiguration.java

@@ -0,0 +1,140 @@
+package com.kxs.common.mq.config;
+
+import cn.hutool.core.util.ClassUtil;
+import cn.hutool.core.util.ReflectUtil;
+import com.kxs.common.core.constant.CommonConstants;
+import com.kxs.common.core.util.SpringContextHolder;
+import com.kxs.common.mq.enums.MQSendTypeEnum;
+import com.kxs.common.mq.model.AbstractMQ;
+import jakarta.annotation.PostConstruct;
+import org.springframework.amqp.core.*;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.beans.factory.support.BeanDefinitionBuilder;
+import org.springframework.stereotype.Component;
+
+import java.util.Objects;
+import java.util.Set;
+
+
+/**
+ * rabbitmq队列配置
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-24
+ */
+@Component
+public class RabbitConfiguration {
+
+	@Autowired
+	private RabbitMQBeanProcessor rabbitMqBeanProcessor;
+
+	public static final String BIND_PREFIX = "BIND_";
+	public static final String DEAD_PREFIX = "DEAD_";
+
+	/**
+	 * 普通交换机
+	 */
+	public static final String EXCHANGE_TOPIC_RANCH ="kxs_direct_ranch";
+
+	/**
+	 * 死信交换机
+	 */
+	public static final String EXCHANGE_DEAD_RANCH ="kxs_dead_ranch";
+
+	/**
+	 * 广播交换机
+	 */
+	public static final String EXCHANGE_FANOUT_RANCH ="kxs_fanout_ranch";
+
+	/**
+	 * 延时交换机
+	 */
+	public static final String EXCHANGE_DELAY_RANCH="kxs_delay_ranch";
+
+	/** 注入延迟交换机Bean **/
+	@Autowired
+	@Qualifier(EXCHANGE_TOPIC_RANCH)
+	private DirectExchange directExchange;
+
+	/** 注入延迟交换机Bean **/
+	@Autowired
+	@Qualifier(EXCHANGE_FANOUT_RANCH)
+	private FanoutExchange fanoutExchange;
+
+	/** 注入延迟交换机Bean **/
+	@Autowired
+	@Qualifier(EXCHANGE_DELAY_RANCH)
+	private CustomExchange delayedExchange;
+
+	/** 死信交换机 **/
+	@Autowired
+	@Qualifier(EXCHANGE_DEAD_RANCH)
+	private DirectExchange deadExchange;
+
+
+	@PostConstruct
+	public void init(){
+
+		// 获取到所有的MQ定义
+		Set<Class<?>> set = ClassUtil.scanPackageBySuper(CommonConstants.PACKAGE_NAME, AbstractMQ.class);
+
+		for (Class<?> aClass : set) {
+			// 实例化
+			AbstractMQ amq = (AbstractMQ) ReflectUtil.newInstance(aClass);
+
+			// 注册正常Queue === new Queue(name),  queue名称/bean名称 = mqName
+//			rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(amq.getQueueName(),
+//					BeanDefinitionBuilder.rootBeanDefinition(Queue.class).addConstructorArgValue(amq.getQueueName()).getBeanDefinition());
+			rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(amq.getQueueName(),
+					BeanDefinitionBuilder.genericBeanDefinition(Queue.class, () ->
+						QueueBuilder.durable(amq.getQueueName()).deadLetterExchange(EXCHANGE_DEAD_RANCH)
+								.deadLetterRoutingKey(DEAD_PREFIX + amq.getQueueName()).build()
+					).getBeanDefinition());
+
+			// 为每个队列 注册死信队列
+			rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(DEAD_PREFIX + amq.getQueueName(),
+					BeanDefinitionBuilder.rootBeanDefinition(Queue.class).addConstructorArgValue(DEAD_PREFIX + amq.getQueueName()).getBeanDefinition());
+
+			//普通队列
+			if(amq.getMqType() == MQSendTypeEnum.QUEUE){
+				rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(BIND_PREFIX + amq.getQueueName(),
+						BeanDefinitionBuilder.genericBeanDefinition(Binding.class, () ->
+								BindingBuilder.bind(SpringContextHolder.getBean(amq.getQueueName(), Queue.class))
+										.to(directExchange).with(amq.getQueueName())
+						).getBeanDefinition()
+				);
+			}
+			// 广播模式
+			if(amq.getMqType() == MQSendTypeEnum.FANOUT){
+				rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(BIND_PREFIX + amq.getQueueName(),
+						BeanDefinitionBuilder.genericBeanDefinition(Binding.class, () ->
+								BindingBuilder.bind(Objects.requireNonNull(SpringContextHolder.getBean(amq.getQueueName(), Queue.class)))
+										.to(fanoutExchange)
+						).getBeanDefinition()
+				);
+			}
+			//延迟队列
+			if(amq.getMqType()== MQSendTypeEnum.DELAY){
+				// 延迟交换机与Queue进行绑定, 绑定Bean名称 = mqName_DelayedBind
+				rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(BIND_PREFIX + amq.getQueueName() + "_delayed",
+						BeanDefinitionBuilder.genericBeanDefinition(Binding.class, () ->
+								BindingBuilder.bind(Objects.requireNonNull(SpringContextHolder.getBean(amq.getQueueName(), Queue.class)))
+										.to(delayedExchange).with(amq.getQueueName()).noargs()
+
+						).getBeanDefinition()
+				);
+			}
+
+			//死信队列绑定死信交换机
+			rabbitMqBeanProcessor.beanDefinitionRegistry.registerBeanDefinition(DEAD_PREFIX + BIND_PREFIX + amq.getQueueName(),
+					BeanDefinitionBuilder.genericBeanDefinition(Binding.class, () ->
+							BindingBuilder.bind(SpringContextHolder.getBean(DEAD_PREFIX + amq.getQueueName(), Queue.class))
+									.to(deadExchange).with(DEAD_PREFIX + amq.getQueueName())
+					).getBeanDefinition()
+			);
+
+		}
+	}
+
+}

+ 79 - 0
kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/config/RabbitMQBeanProcessor.java

@@ -0,0 +1,79 @@
+package com.kxs.common.mq.config;
+
+import org.springframework.amqp.core.*;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
+import org.springframework.beans.factory.support.BeanDefinitionRegistry;
+import org.springframework.beans.factory.support.BeanDefinitionRegistryPostProcessor;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+
+import java.util.HashMap;
+import java.util.Map;
+
+
+/**
+ *   将spring容器的 [bean注册器]放置到属性中,为 RabbitConfig提供访问。
+ *   顺序:
+ *   1. postProcessBeanDefinitionRegistry (存放注册器)
+ *   2. postProcessBeanFactory (没有使用)
+ *   3. 注册延迟消息交换机的bean: delayedExchange
+ *   4. 动态配置RabbitMQ所需的bean。
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+@Configuration
+public class RabbitMQBeanProcessor implements BeanDefinitionRegistryPostProcessor {
+
+    /** bean注册器 **/
+    protected BeanDefinitionRegistry beanDefinitionRegistry;
+
+    @Override
+    public void postProcessBeanDefinitionRegistry(BeanDefinitionRegistry beanDefinitionRegistry) throws BeansException {
+        this.beanDefinitionRegistry = beanDefinitionRegistry;
+    }
+
+    @Override
+    public void postProcessBeanFactory(ConfigurableListableBeanFactory configurableListableBeanFactory) throws BeansException {
+    }
+
+
+
+    /**
+     *  普通交换机
+     */
+    @Bean(name = RabbitConfiguration.EXCHANGE_TOPIC_RANCH)
+    DirectExchange directExchange(){
+
+        return new DirectExchange(RabbitConfiguration.EXCHANGE_TOPIC_RANCH,true,false);
+    }
+
+    /**
+     *  死信交换机
+     */
+    @Bean(name = RabbitConfiguration.EXCHANGE_DEAD_RANCH)
+    DirectExchange deadExchange(){
+
+        return new DirectExchange(RabbitConfiguration.EXCHANGE_DEAD_RANCH,true,false);
+    }
+
+    /**
+     * 广播交换机
+     */
+    @Bean(name = RabbitConfiguration.EXCHANGE_FANOUT_RANCH)
+    FanoutExchange fanoutExchange(){
+        return new FanoutExchange(RabbitConfiguration.EXCHANGE_FANOUT_RANCH, true, false);
+    }
+
+    /**
+     * 延时交换机
+     */
+    @Bean(name = RabbitConfiguration.EXCHANGE_DELAY_RANCH)
+    CustomExchange delayExchange(){
+        Map<String, Object> args = new HashMap<>();
+        args.put("x-delayed-type", "direct");
+        return new CustomExchange(RabbitConfiguration.EXCHANGE_DELAY_RANCH,"x-delayed-message",true, false, args);
+    }
+}

+ 18 - 0
kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/enums/MQSendTypeEnum.java

@@ -0,0 +1,18 @@
+
+package com.kxs.common.mq.enums;
+
+
+/**
+ * 定义MQ消息类型
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+public enum MQSendTypeEnum {
+    /** QUEUE - 点对点 (只有1个消费者可消费。 ActiveMQ的queue模式 ) **/
+    QUEUE,
+    /** 广播模式 **/
+    FANOUT,
+    /** 延时消息 **/
+    DELAY
+}

+ 21 - 0
kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/handlers/IMQSender.java

@@ -0,0 +1,21 @@
+package com.kxs.common.mq.handlers;
+
+
+import com.kxs.common.mq.model.AbstractMQ;
+
+
+/**
+ * MQ 消息发送器 接口定义
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+public interface IMQSender {
+
+    /** 推送MQ消息, 实时 **/
+    void send(AbstractMQ mqModel);
+
+    /** 推送MQ消息, 延迟接收,单位:s **/
+    void send(AbstractMQ mqModel, int delay);
+
+}

+ 49 - 0
kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/handlers/rabbit/RabbitMQSender.java

@@ -0,0 +1,49 @@
+package com.kxs.common.mq.handlers.rabbit;
+
+import com.kxs.common.mq.config.RabbitConfiguration;
+import com.kxs.common.mq.enums.MQSendTypeEnum;
+import com.kxs.common.mq.handlers.IMQSender;
+import com.kxs.common.mq.model.AbstractMQ;
+import lombok.RequiredArgsConstructor;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+import org.springframework.stereotype.Component;
+
+/**
+ * rabbitMQ 消息发送器的实现
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+@Component
+@RequiredArgsConstructor
+public class RabbitMQSender implements IMQSender {
+
+    private final RabbitTemplate rabbitTemplate;
+
+    @Override
+    public void send(AbstractMQ mqModel) {
+
+        if(mqModel.getMqType() == MQSendTypeEnum.QUEUE){
+
+            rabbitTemplate.convertAndSend(mqModel.getQueueName(), mqModel.toMessage());
+        }
+        if(mqModel.getMqType() == MQSendTypeEnum.FANOUT){
+
+            rabbitTemplate.convertAndSend(RabbitConfiguration.EXCHANGE_FANOUT_RANCH + mqModel.getQueueName(), mqModel.getQueueName() + "_ROUTE", mqModel.toMessage());
+        }
+
+    }
+
+    @Override
+    public void send(AbstractMQ mqModel, int delay) {
+
+        if(mqModel.getMqType() == MQSendTypeEnum.DELAY){
+
+            rabbitTemplate.convertAndSend(RabbitConfiguration.EXCHANGE_DELAY_RANCH, mqModel.getQueueName(), mqModel.toMessage(), messagePostProcessor ->{
+                messagePostProcessor.getMessageProperties().setDelay(Math.toIntExact(delay * 1000L));
+                return messagePostProcessor;
+            });
+        }
+    }
+
+}

+ 23 - 0
kxs-common/kxs-common-mq/src/main/java/com/kxs/common/mq/model/AbstractMQ.java

@@ -0,0 +1,23 @@
+package com.kxs.common.mq.model;
+
+import com.kxs.common.mq.enums.MQSendTypeEnum;
+import org.springframework.amqp.core.Message;
+
+/**
+ * 全局顶级 MQ队列配置
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+public abstract class AbstractMQ {
+
+
+    /** MQ名称 **/
+    public abstract String getQueueName();
+
+    /** MQ 类型 **/
+    public abstract MQSendTypeEnum getMqType();
+
+    /** 构造MQ消息体 String类型 **/
+    public abstract Message toMessage();
+}

+ 3 - 1
kxs-common/kxs-common-mq/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports

@@ -1 +1,3 @@
-org.apache.rocketmq.spring.autoconfigure.RocketMQAutoConfiguration
+com.kxs.common.mq.config.RabbitConfiguration
+com.kxs.common.mq.config.RabbitMQBeanProcessor
+com.kxs.common.mq.handlers.rabbit.RabbitMQSender

+ 4 - 0
kxs-product/kxs-product-biz/pom.xml

@@ -79,6 +79,10 @@
             <groupId>cn.hutool</groupId>
             <artifactId>hutool-json</artifactId>
         </dependency>
+        <dependency>
+            <groupId>com.kxs</groupId>
+            <artifactId>kxs-common-mq</artifactId>
+        </dependency>
     </dependencies>
 
     <build>

+ 4 - 0
kxs-system/kxs-system-biz/pom.xml

@@ -75,6 +75,10 @@
             <groupId>com.mysql</groupId>
             <artifactId>mysql-connector-j</artifactId>
         </dependency>
+        <dependency>
+            <groupId>com.kxs</groupId>
+            <artifactId>kxs-common-mq</artifactId>
+        </dependency>
         <!--多数据源-->
         <dependency>
             <groupId>com.baomidou</groupId>

+ 69 - 0
kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/amqp/RabbitOrderQueueListener.java

@@ -0,0 +1,69 @@
+package com.kxs.system.biz.amqp;
+
+import com.kxs.common.mq.config.RabbitConfiguration;
+import com.kxs.system.biz.amqp.rabbit.RabbitOrderQueueMQ;
+import com.kxs.system.biz.service.KxsCampService;
+import com.rabbitmq.client.Channel;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.rabbit.annotation.RabbitHandler;
+import org.springframework.amqp.rabbit.annotation.RabbitListener;
+import org.springframework.stereotype.Component;
+
+import java.io.IOException;
+
+/**
+ * rabbit 队列侦听器
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+@Component
+@Slf4j
+@RequiredArgsConstructor
+public class RabbitOrderQueueListener {
+
+	private final KxsCampService kxsCampService;
+
+	/**
+	 * 监听 hello 队列的处理器
+	 *
+	 * @param message 消息
+	 */
+	@RabbitListener(queues = RabbitOrderQueueMQ.QUEUE_NAME, ackMode = "MANUAL")
+	@RabbitHandler
+	public void onMessage(String msg, Message message, Channel channel){
+		log.info("消费端Payload: " + RabbitOrderQueueMQ.parse(msg));
+        try {
+			RabbitOrderQueueMQ.MsgEntity parse = RabbitOrderQueueMQ.parse(msg);
+			//统计到奖金池
+			kxsCampService.prizePoolIsRefreshed(parse);
+			channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
+        } catch (Exception e) {
+			log.error(e.getMessage(), e);
+			try {
+				channel.basicReject(message.getMessageProperties().getDeliveryTag(), false);
+			} catch (IOException ex) {
+				log.error("进入死信队列失败" + ex.getMessage(), ex);
+			}
+        }
+    }
+
+	/**
+	 * 监听死信 队列的处理器
+	 *
+	 * @param message 消息
+	 */
+	@RabbitListener(queues = RabbitConfiguration.DEAD_PREFIX + RabbitOrderQueueMQ.QUEUE_NAME)
+	@RabbitHandler
+	public void onDeadMessage(String msg, Message message, Channel channel){
+		log.info("死信队列Payload: " + RabbitOrderQueueMQ.parse(msg));
+        try {
+			RabbitOrderQueueMQ.MsgEntity parse = RabbitOrderQueueMQ.parse(msg);
+        } catch (Exception e) {
+			log.error(e.getMessage(), e);
+
+        }
+    }
+}

+ 111 - 0
kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/amqp/rabbit/RabbitOrderQueueMQ.java

@@ -0,0 +1,111 @@
+package com.kxs.system.biz.amqp.rabbit;
+
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONObject;
+import com.kxs.common.mq.enums.MQSendTypeEnum;
+import com.kxs.common.mq.model.AbstractMQ;
+import io.swagger.v3.oas.annotations.media.Schema;
+import lombok.AllArgsConstructor;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.NoArgsConstructor;
+import org.springframework.amqp.core.Message;
+import org.springframework.amqp.core.MessageBuilder;
+import org.springframework.amqp.core.MessageDeliveryMode;
+import org.springframework.amqp.core.MessageProperties;
+
+import java.math.BigDecimal;
+import java.util.UUID;
+
+import static io.lettuce.core.pubsub.PubSubOutput.Type.message;
+
+/**
+ * mq 订单普通队列 配置
+ *
+ * @author 没秃顶的码农
+ * @date 2024-04-25
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@NoArgsConstructor
+@AllArgsConstructor
+public class RabbitOrderQueueMQ extends AbstractMQ {
+
+
+    /**
+     * 订单交易队列
+     */
+    public static final String QUEUE_NAME = "QUEUE_ORDER_DIVISION";
+
+    /**
+     * 内置msg 消息体定义
+     **/
+    private MsgEntity msgEntity;
+
+    /**
+     *  定义Msg消息载体
+     **/
+    @Data
+    public static class MsgEntity {
+
+        /**
+         * 支付订单号
+         **/
+        private String orderId;
+
+        /**
+         * 用户 ID
+         */
+        private Integer userId;
+
+        /**
+         * 商品ID
+         */
+        private Integer goodsId;
+
+        /**
+         * 订单金额
+         */
+        private BigDecimal totalPrice;
+
+    }
+
+    @Override
+    public String getQueueName() {
+
+        return QUEUE_NAME;
+    }
+
+    @Override
+    public MQSendTypeEnum getMqType() {
+
+        return MQSendTypeEnum.QUEUE;
+    }
+
+    @Override
+    public Message toMessage() {
+        String message = JSONObject.toJSONString(msgEntity);
+        // 构建消息体
+        return MessageBuilder.withBody(message.getBytes())
+                .setContentType(MessageProperties.CONTENT_TYPE_TEXT_PLAIN)
+                .setDeliveryMode(MessageDeliveryMode.PERSISTENT)
+                .setMessageId(UUID.randomUUID().toString())
+                .build();
+    }
+
+    /**
+     * 构造发送消息
+     */
+    public static RabbitOrderQueueMQ build(MsgEntity message){
+
+        return new RabbitOrderQueueMQ(message);
+    }
+
+    /**
+     * 解析MQ消息, 一般用于接收MQ消息时
+     */
+    public static MsgEntity parse(String msg){
+        return JSON.parseObject(msg, MsgEntity.class);
+    }
+
+}

+ 4 - 0
kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/mapper/KxsCampMapper.java

@@ -37,5 +37,9 @@ public interface KxsCampMapper extends BaseMapper<KxsCamp> {
      * @return {@link KxsCampUser}
      */
     KxsCampUser selectByUserId(@Param("userId") Integer userId, @Param("now") LocalDateTime now);
+
+    KxsCamp selectUserCamp(@Param("userId") Integer userId, @Param("now") LocalDateTime now);
+
+    KxsCampUser selectUsersCamp(@Param("pids") String[] pids, LocalDateTime now);
 }
 

+ 2 - 1
kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/service/KxsCampService.java

@@ -8,6 +8,7 @@ import com.kxs.system.api.model.KxsCamp;
 import com.kxs.system.api.vo.admin.CampUserExcelVO;
 import com.kxs.system.api.vo.kxsapp.camp.CampGetByIdVO;
 import com.kxs.system.api.vo.kxsapp.camp.CampPageVO;
+import com.kxs.system.biz.amqp.rabbit.RabbitOrderQueueMQ;
 import org.springframework.validation.BindingResult;
 
 import java.util.List;
@@ -64,7 +65,7 @@ public interface KxsCampService extends IService<KxsCamp> {
     /**
      * 奖池已刷新
      */
-    void prizePoolIsRefreshed();
+    void prizePoolIsRefreshed(RabbitOrderQueueMQ.MsgEntity parse);
 
     /**
      * 后台分页 页面

+ 245 - 201
kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/service/impl/KxsCampServiceImpl.java

@@ -8,8 +8,6 @@ import cn.hutool.core.util.NumberUtil;
 import cn.hutool.core.util.StrUtil;
 import com.alibaba.fastjson.JSON;
 import com.alibaba.fastjson.JSONObject;
-import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
-import com.baomidou.mybatisplus.core.conditions.update.UpdateWrapper;
 import com.baomidou.mybatisplus.core.metadata.IPage;
 import com.baomidou.mybatisplus.core.toolkit.Wrappers;
 import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
@@ -23,19 +21,17 @@ import com.kxs.common.core.util.RetOps;
 import com.kxs.common.core.util.SysUtils;
 import com.kxs.common.security.util.SecurityUtils;
 import com.kxs.product.api.feign.RemoteKxsProductService;
-import com.kxs.product.api.model.KxsMachine;
 import com.kxs.product.api.vo.kxsapp.shop.ShopOrderGoodsUserVO;
 import com.kxs.system.api.feign.RemoteOldService;
 import com.kxs.system.api.model.KxsCampUser;
 import com.kxs.system.api.vo.admin.CampUserExcelVO;
 import com.kxs.system.api.vo.kxsapp.camp.CampGetByIdVO;
 import com.kxs.system.api.vo.kxsapp.camp.CampPageVO;
+import com.kxs.system.biz.amqp.rabbit.RabbitOrderQueueMQ;
 import com.kxs.system.biz.constant.enums.CampStatusEnum;
 import com.kxs.system.biz.constant.enums.SysErrorTypeEnum;
 import com.kxs.system.biz.mapper.KxsCampMapper;
 import com.kxs.system.api.model.KxsCamp;
-import com.kxs.system.biz.mapper.KxsCampUserMapper;
-import com.kxs.system.biz.service.KxsCampBonusLogService;
 import com.kxs.system.biz.service.KxsCampService;
 import com.kxs.system.biz.service.KxsCampUserService;
 import com.kxs.user.api.feign.RemoteKxsUserService;
@@ -44,12 +40,10 @@ import com.kxs.user.api.model.KxsUser;
 import com.pig4cloud.plugin.excel.vo.ErrorMessage;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
-import org.springframework.scheduling.annotation.Async;
 import org.springframework.scheduling.annotation.EnableAsync;
 import org.springframework.stereotype.Service;
 import org.springframework.validation.BindingResult;
 
-import javax.json.Json;
 import java.math.BigDecimal;
 import java.time.LocalDateTime;
 import java.util.*;
@@ -76,7 +70,6 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
     private final KxsCampUserService kxsCampUserService;
 
 
-
     @Override
     public R addData(KxsCamp param) {
 
@@ -94,11 +87,11 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
     @Override
     public R updateData(KxsCamp param) {
 
-        if(param.getId() == null){
+        if (param.getId() == null) {
             return R.failed(SysErrorTypeEnum.PARAM_ERROR.getDescription());
         }
         KxsCamp kxsCamp = baseMapper.selectById(param.getId());
-        if(kxsCamp == null){
+        if (kxsCamp == null) {
             return R.failed(SysErrorTypeEnum.NOT_DATA.getDescription());
         }
 
@@ -134,19 +127,19 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
                     .getData()
                     .orElse(null);
             KxsCamp kxsCamp = baseMapper.selectById(excel.getCampId());
-            if(kxsCamp == null){
+            if (kxsCamp == null) {
                 return R.failed(excel.getCampId() + "训练营编号不存在,请重新上传!");
             }
 
             if (Objects.isNull(user)) {
                 errorMsg.add(excel.getUserCode() + "创客未找到!");
-            }else{
-                if(user.getRealStatus() != 1){
+            } else {
+                if (user.getRealStatus() != 1) {
                     errorMsg.add(excel.getUserCode() + "创客未实名!");
                 }
                 //判断创客是否已在其他训练营
                 KxsCampUser kxsCampUser = baseMapper.selectByUserId(user.getId(), now);
-                if(kxsCampUser != null){
+                if (kxsCampUser != null) {
                     errorMsg.add(excel.getUserCode() + "已在训练营!");
                 }
             }
@@ -164,8 +157,7 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
                 campUser.setUsername(user.getUsername());
                 campUser.setUserCode(user.getUserCode());
                 kxsCampUserService.save(campUser);
-            }
-            else {
+            } else {
                 // 数据不合法
                 errorMessageList.add(new ErrorMessage(excel.getLineNum(), errorMsg));
             }
@@ -197,7 +189,7 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
     @Override
     public R getByData(Integer id) {
         KxsCamp kxsCamp = baseMapper.selectById(id);
-        if(kxsCamp == null){
+        if (kxsCamp == null) {
             return R.failed(SysErrorTypeEnum.NOT_DATA.getDescription());
         }
         statusChange(kxsCamp, LocalDateTime.now());
@@ -206,6 +198,60 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
         return R.ok(campGetByIdVO);
     }
 
+    @Override
+    public void prizePoolIsRefreshed(RabbitOrderQueueMQ.MsgEntity parse) {
+        //查询正在进行中的活动
+        LocalDateTime now = LocalDateTime.now();
+        //兑换券下单
+        if(parse.getGoodsId() == 11 || parse.getGoodsId() == 12) {
+
+            //订单奖金
+            BigDecimal amount = NumberUtil.mul(NumberUtil.div(parse.getTotalPrice(), 600), 90);
+            //查询此用户是否参加进行中的训练营
+            KxsCamp kxsCamp = baseMapper.selectUserCamp(parse.getUserId(), now);
+            if (kxsCamp != null) {
+                //给训练营加金额
+                kxsCamp.setBonusPool(NumberUtil.add(kxsCamp.getBonusPool(), amount));
+                kxsCamp.setOrderNum(kxsCamp.getOrderNum() + 1);
+                baseMapper.updateById(kxsCamp);
+
+                //给创客团队加金额
+                KxsCampUser kxsCampUser = kxsCampUserService.getOne(Wrappers.<KxsCampUser>lambdaQuery()
+                        .eq(KxsCampUser::getCampId, kxsCamp.getId())
+                        .eq(KxsCampUser::getUserId, parse.getUserId()));
+                kxsCampUser.setTeamOrderPool(NumberUtil.add(kxsCampUser.getTeamOrderPool(), amount));
+                kxsCampUser.setTeamOrderNum(kxsCampUser.getTeamOrderNum() + 1);
+                kxsCampUserService.updateById(kxsCampUser);
+                log.info("用户{}参加训练营{},订单ID{}, 奖金池{},订单数{}", parse.getUserId(), kxsCamp.getId(), parse.getOrderId(), kxsCamp.getBonusPool(), kxsCamp.getOrderNum());
+                return;
+            }
+            //查找此用户上级倒序后参加了训练营,把奖金加入到自己和训练营里
+            R<KxsUser> kxsUserR = remoteKxsUserService.loadUserById(parse.getUserId(), SecurityConstants.FROM_IN);
+            KxsUser user = RetOps.of(kxsUserR)
+                    .getData()
+                    .orElseThrow(() -> new GlobalCustomerException(ErrorTypeEnum.USER_NOT_FOUND.getDescription()));
+            //下单用户的上级路径
+            String[] pidPaths = user.getPidPath().split(",");
+            //翻转上级从最近的开始查询
+            Arrays.sort(pidPaths, Collections.reverseOrder());
+            KxsCampUser campUser = baseMapper.selectUsersCamp(pidPaths, now);
+            if (campUser != null) {
+                //给训练营加金额
+                KxsCamp camp = baseMapper.selectById(campUser.getCampId());
+                camp.setBonusPool(NumberUtil.add(camp.getBonusPool(), amount));
+                camp.setOrderNum(camp.getOrderNum() + 1);
+                baseMapper.updateById(camp);
+
+                //给创客团队加金额
+                campUser.setTeamOrderPool(NumberUtil.add(campUser.getTeamOrderPool(), amount));
+                campUser.setTeamOrderNum(campUser.getTeamOrderNum() + 1);
+                kxsCampUserService.updateById(campUser);
+                log.info("用户{}的上级{}参加训练营{},订单ID{}, 奖金池{},订单数{}", parse.getUserId(), campUser.getUserId(), camp.getId(), parse.getOrderId(), camp.getBonusPool(), camp.getOrderNum());
+            }
+        }
+
+    }
+
     @Override
     public IPage<KxsCamp> getBySysPage(Page<KxsCamp> page, KxsCamp param) {
         LocalDateTime now = LocalDateTime.now();
@@ -217,208 +263,206 @@ public class KxsCampServiceImpl extends ServiceImpl<KxsCampMapper, KxsCamp> impl
         return pageData;
     }
 
-    @Override
-    @Async("getAsyncExecutor")
-    public void prizePoolIsRefreshed() {
-        //查询正在进行中的活动
-        LocalDateTime now = LocalDateTime.now();
-        //定时任务执行中时,查询当前正在进行的活动,由于定时任务执行时间偏差,当前时间减1一分钟判断
-        List<KxsCamp> list = baseMapper.selectList(Wrappers.<KxsCamp>lambdaQuery().le(KxsCamp::getStartTime, now).ge(KxsCamp::getEndTime, now.minusMinutes(1)));
-        //查询活动参与的创客
-        for (KxsCamp kxsCamp : list) {
-            //初始化奖金池
-            BigDecimal totalBonusPool = BigDecimal.ZERO;
-            //开机数
-            int openNum = 0;
-            //下单数
-            int orderNum = 0;
-
-            List<KxsCampUser> campUsers = kxsCampUserService.list(Wrappers.<KxsCampUser>lambdaQuery().eq(KxsCampUser::getCampId, kxsCamp.getId()));
-            //清除参与创客的所有统计数据然后封装在map里方便统计
-            Map<Integer, KxsCampUser> campUsersMap = new HashMap<>();
-            for (KxsCampUser campUser : campUsers) {
-                KxsCampUser kxsCampUser = new KxsCampUser();
-                BeanUtil.copyProperties(campUser, kxsCampUser);
-                kxsCampUser.setTeamOrderPool(BigDecimal.ZERO);
-                kxsCampUser.setTeamLeaderPool(BigDecimal.ZERO);
-                kxsCampUser.setTeamOrderNum(0);
-                kxsCampUser.setTeamLeaderNum(0);
-                kxsCampUser.setTeamOpenNum(0);
-                campUsersMap.put(campUser.getUserId(), kxsCampUser);
-            }
-
-            //转换为userID集合
-            List<Integer> userIds = campUsers.stream().map(KxsCampUser::getUserId).toList();
-
-
-            /*
-             * 查询当前活动的时间区间所有电签大机的订单
-             */
-            R<List<ShopOrderGoodsUserVO>> ordersR = remoteKxsProductService.getDateBetweenOrders(LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.NORM_DATETIME_PATTERN)
-                    , LocalDateTimeUtil.format(kxsCamp.getEndTime(), DatePattern.NORM_DATETIME_PATTERN), SecurityConstants.FROM_IN);
-            List<ShopOrderGoodsUserVO> orders = RetOps.of(ordersR)
-                    .getData()
-                    .orElse(Collections.emptyList());
-            //循环区间订单查询用户的上级是否包含参与创客的ID
-            for (ShopOrderGoodsUserVO order : orders) {
-                //如果是参与者自己下的单
-                if (userIds.contains(order.getUserId())) {
-                    int currentNum = NumberUtil.div(order.getTotalPrice(), 600).intValue();
-                    //刷新参与创客的团队累计
-                    KxsCampUser kxsCampUser = campUsersMap.get(order.getUserId());
-                    BigDecimal amount = NumberUtil.mul(NumberUtil.div(order.getTotalPrice(), 600), 90);
-                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
-                    //兑换券下单
-                    kxsCampUser.setTeamOrderPool(NumberUtil.add(kxsCampUser.getTeamOrderPool(), amount));
-                    kxsCampUser.setTeamOrderNum(kxsCampUser.getTeamOrderNum() + currentNum);
-//                    if(order.getId() == 27 || order.getId() == 28){
-//                        //大小盟主下单, 奖金池+60,订单数+1
-//                        kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), 60));
-//                        kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
-//                        currentNum = 1;
+//    @Override
+//    public void prizePoolIsRefreshed(RabbitOrderQueueMQ.MsgEntity parse) {
+//        //查询正在进行中的活动
+//        LocalDateTime now = LocalDateTime.now();
+//        //定时任务执行中时,查询当前正在进行的活动,由于定时任务执行时间偏差,当前时间减1一分钟判断
+//        List<KxsCamp> list = baseMapper.selectList(Wrappers.<KxsCamp>lambdaQuery().le(KxsCamp::getStartTime, now).ge(KxsCamp::getEndTime, now.minusMinutes(1)));
+//        //查询活动参与的创客
+//        for (KxsCamp kxsCamp : list) {
+//            //初始化奖金池
+//            BigDecimal totalBonusPool = BigDecimal.ZERO;
+//            //开机数
+//            int openNum = 0;
+//            //下单数
+//            int orderNum = 0;
+//
+//            List<KxsCampUser> campUsers = kxsCampUserService.list(Wrappers.<KxsCampUser>lambdaQuery().eq(KxsCampUser::getCampId, kxsCamp.getId()));
+//            //清除参与创客的所有统计数据然后封装在map里方便统计
+//            Map<Integer, KxsCampUser> campUsersMap = new HashMap<>();
+//            for (KxsCampUser campUser : campUsers) {
+//                KxsCampUser kxsCampUser = new KxsCampUser();
+//                BeanUtil.copyProperties(campUser, kxsCampUser);
+//                kxsCampUser.setTeamOrderPool(BigDecimal.ZERO);
+//                kxsCampUser.setTeamLeaderPool(BigDecimal.ZERO);
+//                kxsCampUser.setTeamOrderNum(0);
+//                kxsCampUser.setTeamLeaderNum(0);
+//                kxsCampUser.setTeamOpenNum(0);
+//                campUsersMap.put(campUser.getUserId(), kxsCampUser);
+//            }
+//
+//            //转换为userID集合
+//            List<Integer> userIds = campUsers.stream().map(KxsCampUser::getUserId).toList();
+//
+//
+//            /*
+//             * 查询当前活动的时间区间所有电签大机的订单
+//             */
+//            R<List<ShopOrderGoodsUserVO>> ordersR = remoteKxsProductService.getDateBetweenOrders(LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.NORM_DATETIME_PATTERN)
+//                    , LocalDateTimeUtil.format(kxsCamp.getEndTime(), DatePattern.NORM_DATETIME_PATTERN), SecurityConstants.FROM_IN);
+//            List<ShopOrderGoodsUserVO> orders = RetOps.of(ordersR)
+//                    .getData()
+//                    .orElse(Collections.emptyList());
+//            //循环区间订单查询用户的上级是否包含参与创客的ID
+//            for (ShopOrderGoodsUserVO order : orders) {
+//                //如果是参与者自己下的单
+//                if (userIds.contains(order.getUserId())) {
+//                    int currentNum = NumberUtil.div(order.getTotalPrice(), 600).intValue();
+//                    //刷新参与创客的团队累计
+//                    KxsCampUser kxsCampUser = campUsersMap.get(order.getUserId());
+//                    BigDecimal amount = NumberUtil.mul(NumberUtil.div(order.getTotalPrice(), 600), 90);
+//                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
+//                    //兑换券下单
+//                    kxsCampUser.setTeamOrderPool(NumberUtil.add(kxsCampUser.getTeamOrderPool(), amount));
+//                    kxsCampUser.setTeamOrderNum(kxsCampUser.getTeamOrderNum() + currentNum);
+////                    if(order.getId() == 27 || order.getId() == 28){
+////                        //大小盟主下单, 奖金池+60,订单数+1
+////                        kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), 60));
+////                        kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
+////                        currentNum = 1;
+////                    }
+//                    orderNum += currentNum;
+//                    continue;
+//                }
+//
+//                R<KxsUser> kxsUserR = remoteKxsUserService.loadUserById(order.getUserId(), SecurityConstants.FROM_IN);
+//                KxsUser user = RetOps.of(kxsUserR)
+//                        .getData()
+//                        .orElseThrow(() -> new GlobalCustomerException(ErrorTypeEnum.USER_NOT_FOUND.getDescription()));
+//                //下单用户的上级路径
+//                String[] pidPaths = user.getPidPath().split(",");
+//                //通过用户的pidPath倒序查找最近所属的上级
+//                Integer userId = SysUtils.checkUserIsChildren(pidPaths, userIds);
+//                if(userId != null){
+//                    int currentNum = NumberUtil.div(order.getTotalPrice(), 600).intValue();
+//                    //刷新参与创客的团队累计
+//                    KxsCampUser kxsCampUser = campUsersMap.get(userId);
+//                    BigDecimal amount = NumberUtil.mul(NumberUtil.div(order.getTotalPrice(), 600), 90);
+//                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
+//
+//                    kxsCampUser.setTeamOrderPool(NumberUtil.add(kxsCampUser.getTeamOrderPool(), amount));
+//                    kxsCampUser.setTeamOrderNum(kxsCampUser.getTeamOrderNum() + currentNum);
+////                    if(order.getId() == 27 || order.getId() == 28){
+////                        //大小盟主下单, 奖金池+60,订单数+1
+////                        kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), 60));
+////                        kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
+////                        currentNum = 1;
+////                    }
+//                    orderNum += currentNum;
+//                }
+//            }
+//
+//            /*
+//             * 查询活动区间开通的大盟主添加到奖金池
+//             */
+//            R<List<KxsLeader>> leadersR = remoteKxsUserService.getDateBetweenLeaders(LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.NORM_DATETIME_PATTERN)
+//                    , LocalDateTimeUtil.format(kxsCamp.getEndTime(), DatePattern.NORM_DATETIME_PATTERN), SecurityConstants.FROM_IN);
+//            List<KxsLeader> leaders = RetOps.of(leadersR)
+//                    .getData()
+//                    .orElse(Collections.emptyList());
+//            for (KxsLeader leader : leaders) {
+//                //如果是参与者自己下的单
+//                if (userIds.contains(leader.getUserId())) {
+//                    BigDecimal amount;
+//                    if(leader.getLeaderType().equals(LeaderTypeEnum.SMALL_LEADER.getType())){
+//                        amount = new BigDecimal("2000");
+//                    }else{
+//                        amount = new BigDecimal("8000");
 //                    }
-                    orderNum += currentNum;
-                    continue;
-                }
-
-                R<KxsUser> kxsUserR = remoteKxsUserService.loadUserById(order.getUserId(), SecurityConstants.FROM_IN);
-                KxsUser user = RetOps.of(kxsUserR)
-                        .getData()
-                        .orElseThrow(() -> new GlobalCustomerException(ErrorTypeEnum.USER_NOT_FOUND.getDescription()));
-                //下单用户的上级路径
-                String[] pidPaths = user.getPidPath().split(",");
-                //通过用户的pidPath倒序查找最近所属的上级
-                Integer userId = SysUtils.checkUserIsChildren(pidPaths, userIds);
-                if(userId != null){
-                    int currentNum = NumberUtil.div(order.getTotalPrice(), 600).intValue();
-                    //刷新参与创客的团队累计
-                    KxsCampUser kxsCampUser = campUsersMap.get(userId);
-                    BigDecimal amount = NumberUtil.mul(NumberUtil.div(order.getTotalPrice(), 600), 90);
-                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
-
-                    kxsCampUser.setTeamOrderPool(NumberUtil.add(kxsCampUser.getTeamOrderPool(), amount));
-                    kxsCampUser.setTeamOrderNum(kxsCampUser.getTeamOrderNum() + currentNum);
-//                    if(order.getId() == 27 || order.getId() == 28){
-//                        //大小盟主下单, 奖金池+60,订单数+1
-//                        kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), 60));
-//                        kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
-//                        currentNum = 1;
+//                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
+//
+//                    //刷新参与创客的团队累计
+//                    KxsCampUser kxsCampUser = campUsersMap.get(leader.getUserId());
+//                    kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), amount));
+//                    kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
+//                    orderNum += 1;
+//                    log.info("参与者{}自己下单盟主,进入到{}奖金池", leader.getUserId(), kxsCamp.getTitle());
+//                    continue;
+//                }
+//                //查询此盟主的上级
+//                R<KxsUser> kxsUserR = remoteKxsUserService.loadUserById(leader.getUserId(), SecurityConstants.FROM_IN);
+//                KxsUser user = RetOps.of(kxsUserR)
+//                        .getData()
+//                        .orElseThrow(() -> new GlobalCustomerException(ErrorTypeEnum.USER_NOT_FOUND.getDescription()));
+//                //下单用户的上级路径
+//                String[] pidPaths = user.getPidPath().split(",");
+//                //通过用户的pidPath倒序查找最近所属的上级
+//                Integer userId = SysUtils.checkUserIsChildren(pidPaths, userIds);
+//                if(userId != null){
+//                    BigDecimal amount;
+//                    if(leader.getLeaderType().equals(LeaderTypeEnum.SMALL_LEADER.getType())){
+//                        amount = new BigDecimal("2000");
+//                    }else{
+//                        amount = new BigDecimal("8000");
 //                    }
-                    orderNum += currentNum;
-                }
-            }
-
-            /*
-             * 查询活动区间开通的大盟主添加到奖金池
-             */
-            R<List<KxsLeader>> leadersR = remoteKxsUserService.getDateBetweenLeaders(LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.NORM_DATETIME_PATTERN)
-                    , LocalDateTimeUtil.format(kxsCamp.getEndTime(), DatePattern.NORM_DATETIME_PATTERN), SecurityConstants.FROM_IN);
-            List<KxsLeader> leaders = RetOps.of(leadersR)
-                    .getData()
-                    .orElse(Collections.emptyList());
-            for (KxsLeader leader : leaders) {
-                //如果是参与者自己下的单
-                if (userIds.contains(leader.getUserId())) {
-                    BigDecimal amount;
-                    if(leader.getLeaderType().equals(LeaderTypeEnum.SMALL_LEADER.getType())){
-                        amount = new BigDecimal("2000");
-                    }else{
-                        amount = new BigDecimal("8000");
-                    }
-                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
-
-                    //刷新参与创客的团队累计
-                    KxsCampUser kxsCampUser = campUsersMap.get(leader.getUserId());
-                    kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), amount));
-                    kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
-                    orderNum += 1;
-                    log.info("参与者{}自己下单盟主,进入到{}奖金池", leader.getUserId(), kxsCamp.getTitle());
-                    continue;
-                }
-                //查询此盟主的上级
-                R<KxsUser> kxsUserR = remoteKxsUserService.loadUserById(leader.getUserId(), SecurityConstants.FROM_IN);
-                KxsUser user = RetOps.of(kxsUserR)
-                        .getData()
-                        .orElseThrow(() -> new GlobalCustomerException(ErrorTypeEnum.USER_NOT_FOUND.getDescription()));
-                //下单用户的上级路径
-                String[] pidPaths = user.getPidPath().split(",");
-                //通过用户的pidPath倒序查找最近所属的上级
-                Integer userId = SysUtils.checkUserIsChildren(pidPaths, userIds);
-                if(userId != null){
-                    BigDecimal amount;
-                    if(leader.getLeaderType().equals(LeaderTypeEnum.SMALL_LEADER.getType())){
-                        amount = new BigDecimal("2000");
-                    }else{
-                        amount = new BigDecimal("8000");
-                    }
-                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
-
-                    //刷新参与创客的团队累计
-                    KxsCampUser kxsCampUser = campUsersMap.get(userId);
-                    kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), amount));
-                    kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
-                    orderNum += 1;
-                    log.info("创客:{}的团队购买盟主:{},进入到:{}奖金池", userId, leader.getUserId(), kxsCamp.getTitle());
-                }
-            }
-
-            /*
-             * 查询活动区间开机数
-             */
-            for (Integer userId : campUsersMap.keySet()) {
-                KxsCampUser kxsCampUser = campUsersMap.get(userId);
-                HashMap<String, Object> data = new HashMap<>();
-                data.put("UserId", userId);
-                data.put("StartTime", LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.PURE_DATE_PATTERN));
-                data.put("EndTime", LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.PURE_DATE_PATTERN));
-                R resR = remoteOldService.teamOpenTotalData(JSON.toJSONString(data));
-                if(resR.getStatus() == 1){
-                    if(resR.getData() != null){
-                        JSONObject jsonObject = JSON.parseObject(JSON.toJSONString(resR.getData()));
-                        int total = jsonObject.getInteger("TeamPosMerchantCount")
-                                + jsonObject.getInteger("TeamSimMerchantCount")
-                                + jsonObject.getInteger("TeamMpMerchantCount");
-                        openNum += total;
-                        kxsCampUser.setTeamOpenNum(kxsCampUser.getTeamOpenNum() + total);
-                    }
-                }
-            }
-
-            kxsCamp.setBonusPool(totalBonusPool);
-            kxsCamp.setOrderNum(orderNum);
-            kxsCamp.setOpenNum(openNum);
-            baseMapper.updateById(kxsCamp);
-            //刷新参与创客的统计
-            kxsCampUserService.updateBatchById(new ArrayList<>(campUsersMap.values()));
-            log.info("训练营{}统计结束", kxsCamp.getCampNum());
-        }
-
-    }
-
+//                    totalBonusPool = NumberUtil.add(totalBonusPool, amount);
+//
+//                    //刷新参与创客的团队累计
+//                    KxsCampUser kxsCampUser = campUsersMap.get(userId);
+//                    kxsCampUser.setTeamLeaderPool(NumberUtil.add(kxsCampUser.getTeamLeaderPool(), amount));
+//                    kxsCampUser.setTeamLeaderNum(kxsCampUser.getTeamLeaderNum() + 1);
+//                    orderNum += 1;
+//                    log.info("创客:{}的团队购买盟主:{},进入到:{}奖金池", userId, leader.getUserId(), kxsCamp.getTitle());
+//                }
+//            }
+//
+//            /*
+//             * 查询活动区间开机数
+//             */
+//            for (Integer userId : campUsersMap.keySet()) {
+//                KxsCampUser kxsCampUser = campUsersMap.get(userId);
+//                HashMap<String, Object> data = new HashMap<>();
+//                data.put("UserId", userId);
+//                data.put("StartTime", LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.PURE_DATE_PATTERN));
+//                data.put("EndTime", LocalDateTimeUtil.format(kxsCamp.getStartTime(), DatePattern.PURE_DATE_PATTERN));
+//                R resR = remoteOldService.teamOpenTotalData(JSON.toJSONString(data));
+//                if(resR.getStatus() == 1){
+//                    if(resR.getData() != null){
+//                        JSONObject jsonObject = JSON.parseObject(JSON.toJSONString(resR.getData()));
+//                        int total = jsonObject.getInteger("TeamPosMerchantCount")
+//                                + jsonObject.getInteger("TeamSimMerchantCount")
+//                                + jsonObject.getInteger("TeamMpMerchantCount");
+//                        openNum += total;
+//                        kxsCampUser.setTeamOpenNum(kxsCampUser.getTeamOpenNum() + total);
+//                    }
+//                }
+//            }
+//
+//            kxsCamp.setBonusPool(totalBonusPool);
+//            kxsCamp.setOrderNum(orderNum);
+//            kxsCamp.setOpenNum(openNum);
+//            baseMapper.updateById(kxsCamp);
+//            //刷新参与创客的统计
+//            kxsCampUserService.updateBatchById(new ArrayList<>(campUsersMap.values()));
+//            log.info("训练营{}统计结束", kxsCamp.getCampNum());
+//        }
+//
+//    }
 
 
     /**
      * 状态变更
      *
      * @param campT 营地
-     * @param now  现在
+     * @param now   现在
      */
-    private <T> void statusChange(T campT, LocalDateTime now){
-        if(campT instanceof KxsCamp camp){
-            if(camp.getEndTime().isBefore(now)){
+    private <T> void statusChange(T campT, LocalDateTime now) {
+        if (campT instanceof KxsCamp camp) {
+            if (camp.getEndTime().isBefore(now)) {
                 camp.setStatus(CampStatusEnum.STATUS_END.getType());
             } else if (LocalDateTimeUtil.isIn(now, camp.getStartTime(), camp.getEndTime(), false, false)) {
                 camp.setStatus(CampStatusEnum.STATUS_NORMAL.getType());
-            }else {
+            } else {
                 camp.setStatus(CampStatusEnum.STATUS_CLOSE.getType());
             }
         }
-        if(campT instanceof CampPageVO camp){
-            if(camp.getEndTime().isBefore(now)){
+        if (campT instanceof CampPageVO camp) {
+            if (camp.getEndTime().isBefore(now)) {
                 camp.setStatus(CampStatusEnum.STATUS_END.getType());
             } else if (LocalDateTimeUtil.isIn(now, camp.getStartTime(), camp.getEndTime(), false, false)) {
                 camp.setStatus(CampStatusEnum.STATUS_NORMAL.getType());
-            }else {
+            } else {
                 camp.setStatus(CampStatusEnum.STATUS_CLOSE.getType());
             }
         }

+ 1 - 1
kxs-system/kxs-system-biz/src/main/java/com/kxs/system/biz/task/KxsSystemTask.java

@@ -44,7 +44,7 @@ public class KxsSystemTask {
     @GetMapping("/prizePoolIsRefreshed")
     public void prizePoolIsRefreshed(@RequestParam("param") String para)  {
 
-        kxsCampService.prizePoolIsRefreshed();
+//        kxsCampService.prizePoolIsRefreshed();
 
     }
 

+ 28 - 0
kxs-system/kxs-system-biz/src/main/resources/mapper/KxsCampMapper.xml

@@ -48,5 +48,33 @@
             and a.end_time <![CDATA[ >= ]]> #{now}
         </where>
     </select>
+    <select id="selectUserCamp" resultType="com.kxs.system.api.model.KxsCamp">
+        SELECT a.id, a.status, a.camp_num, a.camp_type, a.pic_url, a.title, bonus_pool, open_num, order_num
+        from kxs_camp a
+        left join kxs_camp_user b on a.id = b.camp_id
+        <where>
+            and a.del_flag = 0
+            and b.user_id =  #{userId}
+            and a.start_time <![CDATA[ <= ]]> #{now}
+            and a.end_time <![CDATA[ >= ]]> #{now}
+        </where>
+        limit 1
+    </select>
+    <select id="selectUsersCamp" resultType="com.kxs.system.api.model.KxsCampUser">
+        SELECT b.id, b.camp_id, b.user_id, b.team_order_num, b.team_order_pool, b.team_leader_pool, b.team_open_num, b.team_leader_num
+        from kxs_camp a
+        left join kxs_camp_user b on a.id = b.camp_id
+        <where>
+            and a.del_flag = 0
+            and b.user_id in
+            <foreach collection="pids" item="userId" index="index" open="(" close=")" separator=",">
+                #{userId}
+            </foreach>
+            and a.start_time <![CDATA[ <= ]]> #{now}
+            and a.end_time <![CDATA[ >= ]]> #{now}
+        </where>
+        ORDER BY FIELD(id, <foreach collection="pids" item="userId" separator=","> #{userId} </foreach>)
+        limit 1
+    </select>
 
 </mapper>