DtsController.java 8.2 KB

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