DtsController.java 6.6 KB

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