DtsController.java 3.4 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192
  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.KxsDtsLog;
  7. import com.kxs.transfer.api.model.table.DMLData;
  8. import com.kxs.transfer.api.service.KxsDtsLogService;
  9. import com.kxs.transfer.api.service.user.KxsUserAddressService;
  10. import com.kxs.transfer.api.service.user.KxsUserService;
  11. import lombok.RequiredArgsConstructor;
  12. import lombok.extern.slf4j.Slf4j;
  13. import org.springframework.web.bind.annotation.RestController;
  14. import java.util.List;
  15. /**
  16. * <p>
  17. * DTS控制器
  18. * 在此控制器里做数据分发
  19. * 比如将数据分发到不同的数据源中
  20. * </p>
  21. *
  22. * @author 没秃顶的码农
  23. * @date 2024-01-25
  24. */
  25. @RestController
  26. @Slf4j
  27. @RequiredArgsConstructor
  28. public class DtsController {
  29. private final KxsDtsLogService kxsDtsLogService;
  30. private final KxsUserService kxsUserService;
  31. private final KxsUserAddressService kxsUserAddressService;
  32. @DtsMsgListener
  33. public void dtsListener(Long dataId, DMLData dmlData, DefaultUserRecord record){
  34. //过滤无效数据
  35. if("Users".equals(dmlData.getTableName())){
  36. List<String> changeFieldList = dmlData.getChangeFieldList();
  37. //过滤Users老表记录登陆设备的信息日志
  38. if (changeFieldList.size() == 2 && changeFieldList.contains("DeviceId") && changeFieldList.contains("DeviceType")){
  39. record.commit(String.valueOf(record.getSourceTimestamp()));
  40. return;
  41. }
  42. if (changeFieldList.size() == 1 && changeFieldList.contains("DeviceId")){
  43. record.commit(String.valueOf(record.getSourceTimestamp()));
  44. return;
  45. }
  46. //过滤Users本月交易额字段
  47. if (changeFieldList.size() == 1 && changeFieldList.contains("ThisMonthTrade")){
  48. record.commit(String.valueOf(record.getSourceTimestamp()));
  49. return;
  50. }
  51. }
  52. //去重数据
  53. KxsDtsLog dtsLog = kxsDtsLogService.getOne(Wrappers.<KxsDtsLog>lambdaQuery().eq(KxsDtsLog::getDataId, dataId));
  54. if(dtsLog != null){
  55. log.info("dts的数据重复:{},dataID:{}", dmlData.getValidFieldDataMap(), dataId);
  56. record.commit(String.valueOf(record.getSourceTimestamp()));
  57. return;
  58. }
  59. //用户表
  60. if("Users".equals(dmlData.getTableName())){
  61. kxsUserService.changeUser(dmlData);
  62. }
  63. //用户地址表
  64. if("UserAddress".equals(dmlData.getTableName())){
  65. kxsUserAddressService.changeData(dmlData);
  66. }
  67. //密码操作
  68. if("UserMoveInfo".equals(dmlData.getTableName())){
  69. kxsUserService.changeUserPwd(dmlData);
  70. }
  71. //存储历史数据
  72. KxsDtsLog kxsDtsLog = new KxsDtsLog();
  73. kxsDtsLog.setDataId(dataId);
  74. kxsDtsLog.setContent(JSON.toJSONString(dmlData));
  75. kxsDtsLog.setTableName(dmlData.getTableName());
  76. kxsDtsLog.setOperation(dmlData.getOperation().toString());
  77. kxsDtsLogService.save(kxsDtsLog);
  78. log.info("消费dts的数据表信息:{}", dmlData.getValidFieldDataMap());
  79. record.commit(String.valueOf(record.getSourceTimestamp()));
  80. }
  81. }