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.KxsDtsErrorLog;
import com.kxs.transfer.api.model.KxsDtsLog;
import com.kxs.transfer.api.model.table.DMLData;
import com.kxs.transfer.api.service.KxsDtsErrorLogService;
import com.kxs.transfer.api.service.KxsDtsLogService;
import com.kxs.transfer.api.service.product.KxsShopGoodsService;
import com.kxs.transfer.api.service.sys.KxsMorningLogService;
import com.kxs.transfer.api.service.sys.KxsMorningService;
import com.kxs.transfer.api.service.user.*;
import com.kxs.user.api.model.KxsPartner;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang.exception.ExceptionUtils;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
*
* DTS控制器
* 在此控制器里做数据分发
* 比如将数据分发到不同的数据源中
*
*
* @author 没秃顶的码农
* @date 2024-01-25
*/
@RestController
@Slf4j
@RequiredArgsConstructor
public class DtsController {
private final KxsDtsLogService kxsDtsLogService;
private final KxsDtsErrorLogService kxsDtsErrorLogService;
private final KxsUserService kxsUserService;
private final KxsUserPresetLogService kxsUserPresetLogService;
private final KxsUserAddressService kxsUserAddressService;
private final KxsUserAmountService kxsUserAmountService;
private final KxsUserAdvanceService kxsUserAdvanceService;
private final KxsUserAmountLogService kxsUserAmountLogService;
private final KxsUserWithdrawalService kxsUserWithdrawalService;
private final KxsLeaderService kxsLeaderService;
private final KxsLeaderAmountLogService kxsLeaderAmountLogService;
private final KxsPartnerService kxsPartnerService;
private final KxsMorningService kxsMorningService;
//产品模块
private final KxsShopGoodsService kxsShopGoodsService;
@DtsMsgListener
public void dtsListener(Long dataId, DMLData dmlData, DefaultUserRecord record) {
try {
//过滤无效数据
if ("Users".equals(dmlData.getTableName())) {
List 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")) {
// 如果changeFieldList的大小为1,且changeFieldList包含DeviceId,则提交记录
record.commit(String.valueOf(record.getSourceTimestamp()));
return;
}
//过滤Users本月交易额字段
if (changeFieldList.size() == 1 && changeFieldList.contains("ThisMonthTrade")) {
record.commit(String.valueOf(record.getSourceTimestamp()));
return;
}
//过滤Users库存更新字段
if (changeFieldList.size() == 1 && changeFieldList.contains("StoreStock")) {
record.commit(String.valueOf(record.getSourceTimestamp()));
return;
}
}
//去重数据
KxsDtsLog dtsLog = kxsDtsLogService.getOne(Wrappers.lambdaQuery().eq(KxsDtsLog::getDataId, dataId));
if (dtsLog != null) {
log.info("dts的数据重复:{},dataID:{}", dmlData.getValidFieldDataMap(), dataId);
record.commit(String.valueOf(record.getSourceTimestamp()));
return;
}
if (!"UserMoveInfo".equals(dmlData.getTableName())) {
log.info("开始消费:{}表的数据,dataID:{}, 表名:{}", dmlData.getTableName(), dataId, dmlData.getDatabaseName());
}
//用户表
if ("Users".equals(dmlData.getTableName())) {
kxsUserService.changeUser(dmlData);
}
//用户预设职级表
if ("UserRankWhite".equals(dmlData.getTableName())) {
kxsUserPresetLogService.changeUser(dmlData);
}
//密码操作
if ("UserMoveInfo".equals(dmlData.getTableName())) {
kxsUserService.changeUserPwd(dmlData);
//密码操作不存储原始数据
record.commit(String.valueOf(record.getSourceTimestamp()));
return;
}
//用户地址表
if ("UserAddress".equals(dmlData.getTableName())) {
kxsUserAddressService.changeData(dmlData);
}
//用户账户
if ("UserAccount".equals(dmlData.getTableName()) && "KxsProfitServer".equals(dmlData.getDatabaseName())) {
kxsUserAmountService.changeData(dmlData);
}
//用户预扣款表
if ("ToChargeBackRecord".equals(dmlData.getTableName())) {
kxsUserAdvanceService.changeData(dmlData);
}
//用户账户余额日志
if ("UserAccountRecord".equals(dmlData.getTableName())) {
kxsUserAmountLogService.changeData(dmlData);
}
//用户提现申请记录
if ("UserCashRecord".equals(dmlData.getTableName()) && "KxsProfitServer".equals(dmlData.getDatabaseName())) {
kxsUserWithdrawalService.changeData(dmlData);
}
//盟主表
if ("Leaders".equals(dmlData.getTableName())) {
kxsLeaderService.changeData(dmlData);
}
//盟主金额变动记录表
if ("LeaderReserveRecord".equals(dmlData.getTableName())) {
kxsLeaderAmountLogService.changeData(dmlData);
}
//合伙人表
if ("OperatorRankWhite".equals(dmlData.getTableName())) {
kxsPartnerService.changeData(dmlData);
}
//合伙人账户 -> 合伙人表
if ("UserAccount".equals(dmlData.getTableName()) && "KxsOpServer".equals(dmlData.getDatabaseName())) {
kxsPartnerService.changeAccountData(dmlData);
}
//商品表同步
if ("Products".equals(dmlData.getTableName())) {
kxsShopGoodsService.changeData(dmlData);
}
//每日晨会同步
if ("SchoolMorningMeet".equals(dmlData.getTableName())) {
kxsMorningService.changeData(dmlData);
}
if ("SchoolMorningMeetLog".equals(dmlData.getTableName())) {
kxsMorningService.changeLogData(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);
//删除错误记录,不管有没有都删除,因为很快
kxsDtsErrorLogService.removeById(dataId);
record.commit(String.valueOf(record.getSourceTimestamp()));
} catch (Exception e) {
//详细错误信息
String fullStackTrace = ExceptionUtils.getFullStackTrace(e);
//存储错误数据
KxsDtsErrorLog kxsDtsLog = new KxsDtsErrorLog();
kxsDtsLog.setId(dataId);
kxsDtsLog.setContent(JSON.toJSONString(dmlData));
kxsDtsLog.setTableName(dmlData.getTableName());
kxsDtsLog.setOperation(dmlData.getOperation().toString());
kxsDtsLog.setErrorStr(fullStackTrace);
kxsDtsErrorLogService.save(kxsDtsLog);
record.commit(String.valueOf(record.getSourceTimestamp()));
log.error(fullStackTrace);
}
}
}