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.product.api.model.KxsMerchant;
import com.kxs.system.api.model.KxsSchoolStudy;
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.*;
import com.kxs.transfer.api.service.store.KxsMachineRecycleService;
import com.kxs.transfer.api.service.store.KxsWarehouseService;
import com.kxs.transfer.api.service.sys.KxsMorningLogService;
import com.kxs.transfer.api.service.sys.KxsMorningService;
import com.kxs.transfer.api.service.sys.KxsSchoolStudyService;
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 KxsSchoolStudyService kxsSchoolStudyService;
private final KxsMachineService kxsMachineService;
//产品模块
private final KxsShopGoodsService kxsShopGoodsService;
private final KxsMachinePledgeService kxsMachinePledgeService;
private final KxsMachineRatioService kxsMachineRatioService;
private final KxsTicketService kxsTicketService;
private final KxsTicketTransferService kxsTicketTransferService;
private final KxsMerchantService kxsMerchantService;
//仓库模块
private final KxsWarehouseService kxsWarehouseService;
private final KxsMachineRecycleService kxsMachineRecycleService;
@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);
}
//创客学堂
if ("SchoolMakerStudy".equals(dmlData.getTableName())) {
kxsSchoolStudyService.changeData(dmlData);
}
//
// //仓库表
// if ("StoreHouse".equals(dmlData.getTableName())) {
// kxsWarehouseService.changeData(dmlData);
// }
// //机具表
if ("PosMachinesTwo".equals(dmlData.getTableName())) {
kxsMachineService.changeData(dmlData);
}
//商户表
if ("PosMerchantInfo".equals(dmlData.getTableName())) {
kxsMerchantService.changeData(dmlData);
}
//机具回收表
if ("RecycMachineOrder".equals(dmlData.getTableName())) {
kxsMachineRecycleService.changeData(dmlData);
}
if ("RecycMachineOrderPos".equals(dmlData.getTableName())) {
kxsMachineRecycleService.changeInfoData(dmlData);
}
//机具押金调整记录
if ("MerchantDepositSet".equals(dmlData.getTableName())) {
kxsMachinePledgeService.changeData(dmlData);
}
//机具费率调整记录
if ("PosMachinesFeeChangeRecord".equals(dmlData.getTableName())) {
kxsMachineRatioService.changeData(dmlData);
}
//券码表同步
if ("PosCoupons".equals(dmlData.getTableName())) {
kxsTicketService.changeData(dmlData);
}
//券码划拨表同步
if ("PosCouponOrders".equals(dmlData.getTableName())) {
kxsTicketTransferService.changeData(dmlData);
}
//券码划拨明细表同步
if ("PosCouponRecord".equals(dmlData.getTableName())) {
kxsTicketTransferService.changeInfoData(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);
}
}
}