Sfoglia il codice sorgente

股东大盘分红更新数据时添加版本号

mac 2 anni fa
parent
commit
86daa7dc4d

+ 1 - 1
kxs-quartz/src/main/java/com/kxs/daemon/quartz/task/StatBeanTask.java

@@ -31,7 +31,7 @@ public class StatBeanTask {
      */
     @SneakyThrows
     public String prizePoolSynchronization(String para) {
-        remoteKxsStatService.getUserTradeList(para, SecurityConstants.FROM_IN);
+        remoteKxsStatService.shdUserTradeTotalList(para, SecurityConstants.FROM_IN);
         log.info("股东大盘分红积分计算任务:{},输入参数{}", LocalDateTime.now(), para);
         return SkyQuartzEnum.JOB_LOG_STATUS_SUCCESS.getType();
     }

+ 2 - 2
kxs-stat/kxs-stat-api/src/main/java/com/kxs/stat/api/feign/RemoteKxsStatService.java

@@ -25,8 +25,8 @@ public interface RemoteKxsStatService {
 	 * @param from    调用标志
 	 *
 	 */
-	@GetExchange("/stat-job/getUserTradeList")
-    void getUserTradeList(@RequestParam(value = "para", required = false) String para, @RequestHeader(SecurityConstants.FROM) String from);
+	@GetExchange("/stat-job/shdUserTradeTotalList")
+    void shdUserTradeTotalList(@RequestParam(value = "para", required = false) String para, @RequestHeader(SecurityConstants.FROM) String from);
 
 	/**
 	 * 获取交易统计数据

+ 5 - 5
kxs-stat/kxs-stat-biz/src/main/java/com/kxs/stat/biz/mapper/KxsUserTradeMapper.java

@@ -36,15 +36,15 @@ public interface KxsUserTradeMapper extends BaseMapper<KxsUserTrade> {
     /**
      * 创建临时表 下面增删改查
      */
-    void createTemporaryShdTable();
+    void createTemporaryShdTable(String tableName);
 
-    ShdTradeAmtVO getTempShdTableData(Integer userId);
+    ShdTradeAmtVO getTempShdTableData(@Param("tableName") String tableName, @Param("userId") Integer userId);
 
-    void updateTempShdTableData(ShdTradeAmtVO shdTradeAmtVO);
+    void updateTempShdTableData(@Param("tableName") String tableName, @Param("data") ShdTradeAmtVO shdTradeAmtVO);
 
-    Cursor<ShdTradeAmtVO> getTempShdTableList();
+    Cursor<ShdTradeAmtVO> getTempShdTableList(@Param("tableName") String tableName);
 
-    void insertTempShdTableData(ShdTradeAmtVO shdTradeAmtVO);
+    void insertTempShdTableData(@Param("tableName") String tableName, @Param("data") ShdTradeAmtVO shdTradeAmtVO);
 
     BigDecimal getStatAmount(String tableName);
 }

+ 42 - 48
kxs-stat/kxs-stat-biz/src/main/java/com/kxs/stat/biz/service/impl/KxsUserTradeServiceImpl.java

@@ -26,6 +26,7 @@ import org.springframework.transaction.support.TransactionTemplate;
 
 import java.math.BigDecimal;
 import java.time.LocalDate;
+import java.time.LocalDateTime;
 import java.util.ArrayList;
 import java.util.List;
 
@@ -48,6 +49,7 @@ public class KxsUserTradeServiceImpl extends ServiceImpl<KxsUserTradeMapper, Kxs
     private final KxsLkbTradeMapper kxsLkbTradeMapper;
 
     private final static String PREFIX_TABLE_NAME = "kxs_user_trade_";
+    private final static String PREFIX_TABLE_NAME_TEMP = "kxs_user_trade_temp_";
 
     @Override
     @Async
@@ -61,33 +63,36 @@ public class KxsUserTradeServiceImpl extends ServiceImpl<KxsUserTradeMapper, Kxs
             thisMoth = LocalDateTimeUtil.format(LocalDate.now(), DatePattern.SIMPLE_MONTH_PATTERN);
         }
 
+        String version = LocalDateTimeUtil.format(LocalDateTime.now(), DatePattern.PURE_DATETIME_FORMATTER);
+
         String tableName = PREFIX_TABLE_NAME + thisMoth;
+        String tempTableName = PREFIX_TABLE_NAME_TEMP + thisMoth;
 
         int batchSize = 100;
         TransactionTemplate template = new TransactionTemplate(platformTransactionManager);
         template.execute(status -> {
             //创建临时表
-            baseMapper.createTemporaryShdTable();
+            baseMapper.createTemporaryShdTable(tempTableName);
 
             //查询pos交易
             Cursor<ShdTradeAmtVO> kxsUserTrades = baseMapper.getUserTradeList(tableName);
             for (ShdTradeAmtVO kxsUserTrade : kxsUserTrades) {
-                ShdTradeAmtVO shdTradeAmtVO = baseMapper.getTempShdTableData(kxsUserTrade.getUserId());
+                ShdTradeAmtVO shdTradeAmtVO = baseMapper.getTempShdTableData(tempTableName, kxsUserTrade.getUserId());
                 if(shdTradeAmtVO == null){
                     shdTradeAmtVO = new ShdTradeAmtVO();
                     shdTradeAmtVO.setUserId(kxsUserTrade.getUserId());
                     shdTradeAmtVO.setPosAmt(kxsUserTrade.getPosAmt());
-                    baseMapper.insertTempShdTableData(shdTradeAmtVO);
+                    baseMapper.insertTempShdTableData(tempTableName, shdTradeAmtVO);
                 }else{
                     shdTradeAmtVO.setPosAmt(shdTradeAmtVO.getPosAmt());
-                    baseMapper.updateTempShdTableData(shdTradeAmtVO);
+                    baseMapper.updateTempShdTableData(tempTableName, shdTradeAmtVO);
                 }
 
             }
             //查询lkb交易
             Cursor<KxsLkbTrade> lkbTrades = kxsLkbTradeMapper.getUserTradeMonthList(thisMoth);
             for (KxsLkbTrade lkbTrade : lkbTrades) {
-                ShdTradeAmtVO shdTradeAmtVO = baseMapper.getTempShdTableData(lkbTrade.getUserId());
+                ShdTradeAmtVO shdTradeAmtVO = baseMapper.getTempShdTableData(tempTableName, lkbTrade.getUserId());
                 if(shdTradeAmtVO == null){
                     shdTradeAmtVO = new ShdTradeAmtVO();
                     shdTradeAmtVO.setUserId(lkbTrade.getUserId());
@@ -96,45 +101,67 @@ public class KxsUserTradeServiceImpl extends ServiceImpl<KxsUserTradeMapper, Kxs
                     if(lkbTrade.getActTradeAmt() != null){
                         shdTradeAmtVO.setLkbActAmt(lkbTrade.getActTradeAmt().multiply(BigDecimal.valueOf(4)));
                     }
-                    baseMapper.insertTempShdTableData(shdTradeAmtVO);
+                    baseMapper.insertTempShdTableData(tempTableName, shdTradeAmtVO);
                 }else{
                     shdTradeAmtVO.setLkbAmt(lkbTrade.getTradeAmt());
                     //活动交易 * 4加入
                     if(lkbTrade.getActTradeAmt() != null){
                         shdTradeAmtVO.setLkbActAmt(lkbTrade.getActTradeAmt().multiply(BigDecimal.valueOf(4)));
                     }
-                    baseMapper.updateTempShdTableData(shdTradeAmtVO);
+                    baseMapper.updateTempShdTableData(tempTableName, shdTradeAmtVO);
                 }
             }
             //查询广电交易数据
             Cursor<KxsUserActTrade> kxsUserActTrades = kxsUserActTradeMapper.getUserTradeMonthList(thisMoth);
             for (KxsUserActTrade kxsUserActTrade : kxsUserActTrades) {
-                ShdTradeAmtVO shdTradeAmtVO = baseMapper.getTempShdTableData(kxsUserActTrade.getUserId());
+                ShdTradeAmtVO shdTradeAmtVO = baseMapper.getTempShdTableData(tempTableName, kxsUserActTrade.getUserId());
                 if(shdTradeAmtVO == null){
                     shdTradeAmtVO = new ShdTradeAmtVO();
                     shdTradeAmtVO.setUserId(kxsUserActTrade.getUserId());
                     //广电卡 * 10000加入
                     shdTradeAmtVO.setGdAmt(NumberUtil.mul(kxsUserActTrade.getActNum(), BigDecimal.valueOf(10000)));
-                    baseMapper.insertTempShdTableData(shdTradeAmtVO);
+                    baseMapper.insertTempShdTableData(tempTableName, shdTradeAmtVO);
                 }else{
                     //广电卡 * 10000加入
                     shdTradeAmtVO.setGdAmt(NumberUtil.mul(kxsUserActTrade.getActNum(), BigDecimal.valueOf(10000)));
-                    baseMapper.updateTempShdTableData(shdTradeAmtVO);
+                    baseMapper.updateTempShdTableData(tempTableName, shdTradeAmtVO);
                 }
             }
 
-            Cursor<ShdTradeAmtVO> shdTradeAmtVOList = baseMapper.getTempShdTableList();
-            List<ShdTradeAmtVO> batch = new ArrayList<>();
+            Cursor<ShdTradeAmtVO> shdTradeAmtVOList = baseMapper.getTempShdTableList(tempTableName);
+            List<IntegralStatDTO> batch = new ArrayList<>();
             for (ShdTradeAmtVO kxsUserTrade : shdTradeAmtVOList) {
-                batch.add(kxsUserTrade);
+                //pos交易额 + 来客吧非活动交易额
+                BigDecimal totalAmount = BigDecimal.ZERO;
+                if(kxsUserTrade.getPosAmt() != null){
+                    totalAmount = totalAmount.add(kxsUserTrade.getPosAmt());
+                }
+                if(kxsUserTrade.getLkbAmt() != null){
+                    totalAmount = totalAmount.add(kxsUserTrade.getLkbAmt());
+                }
+                if(kxsUserTrade.getGdAmt() != null){
+                    totalAmount = totalAmount.add(kxsUserTrade.getGdAmt());
+                }
+                if(kxsUserTrade.getLkbActAmt() != null){
+                    totalAmount = totalAmount.add(kxsUserTrade.getLkbActAmt());
+                }
+                if(totalAmount.compareTo(BigDecimal.valueOf(3000000)) >= 0){
+                    IntegralStatDTO integralStatDTO = new IntegralStatDTO();
+                    integralStatDTO.setTradeMonth(thisMoth);
+                    integralStatDTO.setUserId(Long.valueOf(kxsUserTrade.getUserId()));
+                    integralStatDTO.setTotalAmount(totalAmount);
+                    integralStatDTO.setVersion(Long.valueOf(version));
+                    batch.add(integralStatDTO);
+                    log.info("用户{}交易额达标:{}", kxsUserTrade.getUserId(), totalAmount);
+                }
                 if (batch.size() >= batchSize) {
-                    remoteUserShdStat(batch, thisMoth);
+                    remoteKxsUserService.integralStat(batch, SecurityConstants.FROM_IN);
                     batch.clear();
                 }
             }
             // 处理最后一批数据
             if (!batch.isEmpty()) {
-                remoteUserShdStat(batch, thisMoth);
+                remoteKxsUserService.integralStat(batch, SecurityConstants.FROM_IN);
             }
             return null;
         });
@@ -147,38 +174,5 @@ public class KxsUserTradeServiceImpl extends ServiceImpl<KxsUserTradeMapper, Kxs
         return baseMapper.getStatAmount(tableName);
     }
 
-
-    public void remoteUserShdStat(List<ShdTradeAmtVO> userTradeList, String month) {
-        List<IntegralStatDTO> params = new ArrayList<>();
-        for (ShdTradeAmtVO datum : userTradeList) {
-
-
-            Integer userId = datum.getUserId();
-
-            //pos交易额 + 来客吧非活动交易额
-            BigDecimal totalAmount = BigDecimal.ZERO;
-            if(datum.getPosAmt() != null){
-                totalAmount = totalAmount.add(datum.getPosAmt());
-            }
-            if(datum.getLkbAmt() != null){
-                totalAmount = totalAmount.add(datum.getLkbAmt());
-            }
-            if(datum.getGdAmt() != null){
-                totalAmount = totalAmount.add(datum.getGdAmt());
-            }
-            if(datum.getLkbActAmt() != null){
-                totalAmount = totalAmount.add(datum.getLkbActAmt());
-            }
-            if(totalAmount.compareTo(BigDecimal.valueOf(3000000)) >= 0){
-                params.add(IntegralStatDTO.builder().userId(Long.valueOf(userId)).tradeMonth(month).totalAmount(totalAmount).build());
-                log.info("用户{}交易额达标:{}", userId, totalAmount);
-            }
-        }
-        remoteKxsUserService.integralStat(params, SecurityConstants.FROM_IN);
-    }
-
-
-
-
 }
 

+ 3 - 3
kxs-stat/kxs-stat-biz/src/main/java/com/kxs/stat/biz/task/KxsStatTaskJob.java

@@ -42,9 +42,9 @@ public class KxsStatTaskJob {
     /**
      * 股东大盘分红积分计算任务
      */
-    @Inner
-    @GetMapping("/getUserTradeList")
-    public void getUserTradeList(@RequestParam(value = "para", required = false) String para) {
+    @Inner(value = false)
+    @GetMapping("/shdUserTradeTotalList")
+    public void shdUserTradeTotalList(@RequestParam(value = "para", required = false) String para) {
         kxsUserTradeService.getUserTradeList(para);
     }
 

+ 31 - 31
kxs-stat/kxs-stat-biz/src/main/resources/mapper/KxsUserTradeMapper.xml

@@ -51,81 +51,81 @@
 
     <!--    临时表操作 下面增删改查-->
     <update id="createTemporaryShdTable">
-        CREATE
-        TEMPORARY TABLE kxs_user_trade_temp (
+        CREATE TEMPORARY TABLE IF NOT EXISTS `${tableName}` (
             user_id INT PRIMARY KEY,
             pos_amt NUMERIC(18,2) default 0,
             lkb_amt NUMERIC(18,2) default 0,
             lkb_act_amt NUMERIC(18,2) default 0,
             gd_amt NUMERIC(18,2) default 0
         );
+        TRUNCATE TABLE `${tableName}`;
     </update>
     <select id="getTempShdTableData" resultType="com.kxs.stat.api.vo.ShdTradeAmtVO">
-        select * from kxs_user_trade_temp
+        select * from `${tableName}`
         <where>
             and user_id = #{userId}
         </where>
     </select>
     <select id="getTempShdTableList" resultType="com.kxs.stat.api.vo.ShdTradeAmtVO" fetchSize="1000">
-        select * from kxs_user_trade_temp
+        select * from `${tableName}`
     </select>
     <select id="getStatAmount" resultType="java.math.BigDecimal">
-        select sum(pro_debit_trade_amt + pro_direct_trade_amt) from ${tableName}
+        select sum(pro_debit_trade_amt + pro_direct_trade_amt) from `${tableName}`
     </select>
     <insert id="insertTempShdTableData">
-        INSERT INTO kxs_user_trade_temp
+        INSERT INTO `${tableName}`
         <trim prefix="(" suffix=")" suffixOverrides=",">
-            <if test="userId != null">
+            <if test="data.userId != null">
                 user_id,
             </if>
-            <if test="posAmt != null">
+            <if test="data.posAmt != null">
                 pos_amt,
             </if>
-            <if test="lkbAmt != null">
+            <if test="data.lkbAmt != null">
                 lkb_amt,
             </if>
-            <if test="lkbActAmt != null">
+            <if test="data.lkbActAmt != null">
                 lkb_act_amt,
             </if>
-            <if test="gdAmt != null">
+            <if test="data.gdAmt != null">
                 gd_amt,
             </if>
         </trim>
         <trim prefix="values(" suffix=")" suffixOverrides=",">
-            <if test="userId != null">
-                #{userId},
+            <if test="data.userId != null">
+                #{data.userId},
             </if>
-            <if test="posAmt != null">
-                #{posAmt},
+            <if test="data.posAmt != null">
+                #{data.posAmt},
             </if>
-            <if test="lkbAmt != null">
-                #{lkbAmt},
+            <if test="data.lkbAmt != null">
+                #{data.lkbAmt},
             </if>
-            <if test="lkbActAmt != null">
-                #{lkbActAmt},
+            <if test="data.lkbActAmt != null">
+                #{data.lkbActAmt},
             </if>
-            <if test="gdAmt != null">
-                #{gdAmt},
+            <if test="data.gdAmt != null">
+                #{data.gdAmt},
             </if>
         </trim>
     </insert>
     <update id="updateTempShdTableData">
-        UPDATE kxs_user_trade_temp
+        UPDATE `${tableName}`
         <set>
-            <if test="posAmt != null">
-                pos_amt = #{posAmt},
+            <if test="data.posAmt != null">
+                pos_amt = #{data.posAmt},
             </if>
-            <if test="lkbAmt != null">
-                lkb_amt = #{lkbAmt},
+            <if test="data.lkbAmt != null">
+                lkb_amt = #{data.lkbAmt},
             </if>
-            <if test="lkbActAmt != null">
-                lkb_act_amt = #{lkbActAmt},
+            <if test="data.lkbActAmt != null">
+                lkb_act_amt = #{data.lkbActAmt},
             </if>
-            <if test="gdAmt != null">
-                gd_amt = #{gdAmt},
+            <if test="data.gdAmt != null">
+                gd_amt = #{data.gdAmt},
             </if>
         </set>
-        where user_id = #{userId}
+        where user_id = #{data.userId}
     </update>
 
 

+ 6 - 2
kxs-user/kxs-user-api/src/main/java/com/kxs/user/api/dto/kxsapp/IntegralStatDTO.java

@@ -2,8 +2,8 @@ package com.kxs.user.api.dto.kxsapp;
 
 import lombok.Builder;
 import lombok.Data;
+import lombok.NoArgsConstructor;
 
-import java.io.Serial;
 import java.io.Serializable;
 import java.math.BigDecimal;
 
@@ -14,7 +14,6 @@ import java.math.BigDecimal;
  * @date 2024-05-09
  */
 @Data
-@Builder
 public class IntegralStatDTO implements Serializable {
 
     /**
@@ -31,4 +30,9 @@ public class IntegralStatDTO implements Serializable {
      * 月份
      */
     private String tradeMonth;
+
+    /**
+     * 更新版本
+     */
+    private Long version;
 }

+ 1 - 1
kxs-user/kxs-user-api/src/main/java/com/kxs/user/api/feign/RemoteKxsUserService.java

@@ -159,7 +159,7 @@ public interface RemoteKxsUserService {
 	 *
 	 * @return {@link R}
 	 */
-	@GetExchange("/shareholder/integralStat")
+	@PostExchange("/shareholder/integralStat")
 	void integralStat(@RequestBody List<IntegralStatDTO> integralStatDTOS,@RequestHeader(SecurityConstants.FROM) String from);
 
 	/**

+ 1 - 1
kxs-user/kxs-user-biz/src/main/java/com/kxs/user/biz/controller/kxsapp/ShareholderController.java

@@ -84,7 +84,7 @@ public class ShareholderController {
      * 每日计算
      */
     @Inner
-    @GetMapping("/integralStat")
+    @PostMapping("/integralStat")
     public void integralStat(@RequestBody List<IntegralStatDTO> integralStatDTOS) {
         kxsShdScoreService.integralStat(integralStatDTOS);
     }