DtsController.java 7.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. package com.kxs.transfer.api.controller;
  2. import com.alibaba.fastjson.JSON;
  3. import com.aliyun.dts.subscribe.clients.record.DefaultUserRecord;
  4. import com.baomidou.mybatisplus.core.toolkit.Wrappers;
  5. import com.kxs.transfer.api.annotation.DtsMsgListener;
  6. import com.kxs.transfer.api.model.KxsDtsErrorLog;
  7. import com.kxs.transfer.api.model.KxsDtsLog;
  8. import com.kxs.transfer.api.model.table.DMLData;
  9. import com.kxs.transfer.api.service.KxsDtsErrorLogService;
  10. import com.kxs.transfer.api.service.KxsDtsLogService;
  11. import com.kxs.transfer.api.service.product.KxsShopGoodsService;
  12. import com.kxs.transfer.api.service.user.*;
  13. import com.kxs.user.api.model.KxsPartner;
  14. import lombok.RequiredArgsConstructor;
  15. import lombok.extern.slf4j.Slf4j;
  16. import org.apache.commons.lang.exception.ExceptionUtils;
  17. import org.springframework.web.bind.annotation.RestController;
  18. import java.util.List;
  19. /**
  20. * <p>
  21. * DTS控制器
  22. * 在此控制器里做数据分发
  23. * 比如将数据分发到不同的数据源中
  24. * </p>
  25. *
  26. * @author 没秃顶的码农
  27. * @date 2024-01-25
  28. */
  29. @RestController
  30. @Slf4j
  31. @RequiredArgsConstructor
  32. public class DtsController {
  33. private final KxsDtsLogService kxsDtsLogService;
  34. private final KxsDtsErrorLogService kxsDtsErrorLogService;
  35. private final KxsUserService kxsUserService;
  36. private final KxsUserPresetLogService kxsUserPresetLogService;
  37. private final KxsUserAddressService kxsUserAddressService;
  38. private final KxsUserAmountService kxsUserAmountService;
  39. private final KxsUserAdvanceService kxsUserAdvanceService ;
  40. private final KxsUserAmountLogService kxsUserAmountLogService;
  41. private final KxsUserWithdrawalService kxsUserWithdrawalService;
  42. private final KxsLeaderService kxsLeaderService;
  43. private final KxsLeaderAmountLogService kxsLeaderAmountLogService;
  44. private final KxsPartnerService kxsPartnerService;
  45. //产品模块
  46. private final KxsShopGoodsService kxsShopGoodsService;
  47. @DtsMsgListener
  48. public void dtsListener(Long dataId, DMLData dmlData, DefaultUserRecord record) {
  49. try {
  50. //过滤无效数据
  51. if ("Users".equals(dmlData.getTableName())) {
  52. List<String> changeFieldList = dmlData.getChangeFieldList();
  53. //过滤Users老表记录登陆设备的信息日志
  54. if (changeFieldList.size() == 2 && changeFieldList.contains("DeviceId") && changeFieldList.contains("DeviceType")) {
  55. record.commit(String.valueOf(record.getSourceTimestamp()));
  56. return;
  57. }
  58. if (changeFieldList.size() == 1 && changeFieldList.contains("DeviceId")) {
  59. // 如果changeFieldList的大小为1,且changeFieldList包含DeviceId,则提交记录
  60. record.commit(String.valueOf(record.getSourceTimestamp()));
  61. return;
  62. }
  63. //过滤Users本月交易额字段
  64. if (changeFieldList.size() == 1 && changeFieldList.contains("ThisMonthTrade")) {
  65. record.commit(String.valueOf(record.getSourceTimestamp()));
  66. return;
  67. }
  68. //过滤Users库存更新字段
  69. if (changeFieldList.size() == 1 && changeFieldList.contains("StoreStock")) {
  70. record.commit(String.valueOf(record.getSourceTimestamp()));
  71. return;
  72. }
  73. }
  74. //去重数据
  75. KxsDtsLog dtsLog = kxsDtsLogService.getOne(Wrappers.<KxsDtsLog>lambdaQuery().eq(KxsDtsLog::getDataId, dataId));
  76. if (dtsLog != null) {
  77. log.info("dts的数据重复:{},dataID:{}", dmlData.getValidFieldDataMap(), dataId);
  78. record.commit(String.valueOf(record.getSourceTimestamp()));
  79. return;
  80. }
  81. if (!"UserMoveInfo".equals(dmlData.getTableName())) {
  82. log.info("开始消费:{}表的数据,dataID:{}", dmlData.getTableName(), dataId);
  83. }
  84. //用户表
  85. if ("Users".equals(dmlData.getTableName())) {
  86. kxsUserService.changeUser(dmlData);
  87. }
  88. //用户预设职级表
  89. if ("UserRankWhite".equals(dmlData.getTableName())) {
  90. kxsUserPresetLogService.changeUser(dmlData);
  91. }
  92. //密码操作
  93. if ("UserMoveInfo".equals(dmlData.getTableName())) {
  94. kxsUserService.changeUserPwd(dmlData);
  95. //密码操作不存储原始数据
  96. record.commit(String.valueOf(record.getSourceTimestamp()));
  97. return;
  98. }
  99. //用户地址表
  100. if ("UserAddress".equals(dmlData.getTableName())) {
  101. kxsUserAddressService.changeData(dmlData);
  102. }
  103. //用户账户
  104. if ("UserAccount".equals(dmlData.getTableName()) && "KxsProfitServer".equals(dmlData.getDatabaseName())) {
  105. kxsUserAmountService.changeData(dmlData);
  106. }
  107. //用户预扣款表
  108. if ("ToChargeBackRecord".equals(dmlData.getTableName())) {
  109. kxsUserAdvanceService.changeData(dmlData);
  110. }
  111. //用户账户余额日志
  112. if ("UserAccountRecord".equals(dmlData.getTableName())) {
  113. kxsUserAmountLogService.changeData(dmlData);
  114. }
  115. //用户提现申请记录
  116. if ("UserCashRecord".equals(dmlData.getTableName()) && "KxsProfitServer".equals(dmlData.getDatabaseName())) {
  117. kxsUserWithdrawalService.changeData(dmlData);
  118. }
  119. //盟主表
  120. if ("Leaders".equals(dmlData.getTableName()) ) {
  121. kxsLeaderService.changeData(dmlData);
  122. }
  123. //盟主金额变动记录表
  124. if ("LeaderReserveRecord".equals(dmlData.getTableName()) ) {
  125. kxsLeaderAmountLogService.changeData(dmlData);
  126. }
  127. //合伙人表
  128. if ("OperatorRankWhite".equals(dmlData.getTableName()) ) {
  129. kxsPartnerService.changeData(dmlData);
  130. }
  131. //合伙人账户 -> 合伙人表
  132. if ("UserAccount".equals(dmlData.getTableName()) && "KxsOpServer".equals(dmlData.getDatabaseName())) {
  133. kxsPartnerService.changeAccountData(dmlData);
  134. }
  135. //商品表同步
  136. if ("Products".equals(dmlData.getTableName()) ) {
  137. kxsShopGoodsService.changeData(dmlData);
  138. }
  139. //存储历史数据
  140. KxsDtsLog kxsDtsLog = new KxsDtsLog();
  141. kxsDtsLog.setDataId(dataId);
  142. kxsDtsLog.setContent(JSON.toJSONString(dmlData));
  143. kxsDtsLog.setTableName(dmlData.getTableName());
  144. kxsDtsLog.setOperation(dmlData.getOperation().toString());
  145. kxsDtsLogService.save(kxsDtsLog);
  146. //删除错误记录,不管有没有都删除,因为很快
  147. kxsDtsErrorLogService.removeById(dataId);
  148. record.commit(String.valueOf(record.getSourceTimestamp()));
  149. } catch (Exception e) {
  150. //详细错误信息
  151. String fullStackTrace = ExceptionUtils.getFullStackTrace(e);
  152. //存储错误数据
  153. KxsDtsErrorLog kxsDtsLog = new KxsDtsErrorLog();
  154. kxsDtsLog.setId(dataId);
  155. kxsDtsLog.setContent(JSON.toJSONString(dmlData));
  156. kxsDtsLog.setTableName(dmlData.getTableName());
  157. kxsDtsLog.setOperation(dmlData.getOperation().toString());
  158. kxsDtsLog.setErrorStr(fullStackTrace);
  159. kxsDtsErrorLogService.save(kxsDtsLog);
  160. record.commit(String.valueOf(record.getSourceTimestamp()));
  161. log.error(fullStackTrace);
  162. }
  163. }
  164. }