| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253 |
- 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;
- /**
- * <p>
- * DTS控制器
- * 在此控制器里做数据分发
- * 比如将数据分发到不同的数据源中
- * </p>
- *
- * @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<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")) {
- // 如果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.<KxsDtsLog>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);
- }
- }
- }
|