| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192 |
- package com.kxs.transfer.api.controller;
- import com.alibaba.fastjson.JSON;
- import com.aliyun.dts.subscribe.clients.record.DefaultUserRecord;
- import com.baomidou.mybatisplus.core.toolkit.Wrappers;
- import com.kxs.transfer.api.annotation.DtsMsgListener;
- import com.kxs.transfer.api.model.KxsDtsLog;
- import com.kxs.transfer.api.model.table.DMLData;
- import com.kxs.transfer.api.service.KxsDtsLogService;
- import com.kxs.transfer.api.service.user.KxsUserAddressService;
- import com.kxs.transfer.api.service.user.KxsUserService;
- import lombok.RequiredArgsConstructor;
- import lombok.extern.slf4j.Slf4j;
- import org.springframework.web.bind.annotation.RestController;
- import java.util.List;
- /**
- * <p>
- * DTS控制器
- * 在此控制器里做数据分发
- * 比如将数据分发到不同的数据源中
- * </p>
- *
- * @author 没秃顶的码农
- * @date 2024-01-25
- */
- @RestController
- @Slf4j
- @RequiredArgsConstructor
- public class DtsController {
- private final KxsDtsLogService kxsDtsLogService;
- private final KxsUserService kxsUserService;
- private final KxsUserAddressService kxsUserAddressService;
- @DtsMsgListener
- public void dtsListener(Long dataId, DMLData dmlData, DefaultUserRecord record){
- //过滤无效数据
- if("Users".equals(dmlData.getTableName())){
- List<String> changeFieldList = dmlData.getChangeFieldList();
- //过滤Users老表记录登陆设备的信息日志
- if (changeFieldList.size() == 2 && changeFieldList.contains("DeviceId") && changeFieldList.contains("DeviceType")){
- record.commit(String.valueOf(record.getSourceTimestamp()));
- return;
- }
- if (changeFieldList.size() == 1 && changeFieldList.contains("DeviceId")){
- record.commit(String.valueOf(record.getSourceTimestamp()));
- return;
- }
- //过滤Users本月交易额字段
- if (changeFieldList.size() == 1 && changeFieldList.contains("ThisMonthTrade")){
- record.commit(String.valueOf(record.getSourceTimestamp()));
- return;
- }
- }
- //去重数据
- KxsDtsLog dtsLog = kxsDtsLogService.getOne(Wrappers.<KxsDtsLog>lambdaQuery().eq(KxsDtsLog::getDataId, dataId));
- if(dtsLog != null){
- log.info("dts的数据重复:{},dataID:{}", dmlData.getValidFieldDataMap(), dataId);
- record.commit(String.valueOf(record.getSourceTimestamp()));
- return;
- }
- //用户表
- if("Users".equals(dmlData.getTableName())){
- kxsUserService.changeUser(dmlData);
- }
- //用户地址表
- if("UserAddress".equals(dmlData.getTableName())){
- kxsUserAddressService.changeData(dmlData);
- }
- //密码操作
- if("UserMoveInfo".equals(dmlData.getTableName())){
- kxsUserService.changeUserPwd(dmlData);
- }
- //存储历史数据
- KxsDtsLog kxsDtsLog = new KxsDtsLog();
- kxsDtsLog.setDataId(dataId);
- kxsDtsLog.setContent(JSON.toJSONString(dmlData));
- kxsDtsLog.setTableName(dmlData.getTableName());
- kxsDtsLog.setOperation(dmlData.getOperation().toString());
- kxsDtsLogService.save(kxsDtsLog);
- log.info("消费dts的数据表信息:{}", dmlData.getValidFieldDataMap());
- record.commit(String.valueOf(record.getSourceTimestamp()));
- }
- }
|