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.store.api.model.KxsMachineAdvance; import com.kxs.store.api.model.KxsMachineApply; import com.kxs.store.api.model.KxsWarehouseLimitLog; import com.kxs.store.api.model.KxsWarehouseStockLog; import com.kxs.system.api.model.KxsSchoolStudy; import com.kxs.transfer.api.annotation.DtsMsgListener; import com.kxs.transfer.api.mapper.user.KxsUserByStageMapper; 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.product.impl.KxsMachineTrackService; import com.kxs.transfer.api.service.product.impl.KxsShopOrderService; import com.kxs.transfer.api.service.stat.*; import com.kxs.transfer.api.service.store.KxsMachineRecycleService; import com.kxs.transfer.api.service.store.KxsWarehouseService; import com.kxs.transfer.api.service.store.impl.*; import com.kxs.transfer.api.service.sys.*; import com.kxs.transfer.api.service.user.*; import com.kxs.transfer.api.service.user.impl.KxsUserPresetBeforeService; import com.kxs.transfer.api.service.user.impl.KxsUserReRealService; import com.kxs.user.api.model.KxsLeaderAmountLog; import com.kxs.user.api.model.KxsPartner; import com.kxs.user.api.model.KxsUserByStage; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang.exception.ExceptionUtils; import org.springframework.beans.factory.annotation.Value; 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 KxsPartnerService kxsPartnerService; private final KxsMorningService kxsMorningService; private final KxsSchoolStudyService kxsSchoolStudyService; private final KxsSysMsgService kxsSysMsgService; private final KxsBankChangeLogService kxsBankChangeLogService; private final KxsUserByStageService kxsUserByStageService; private final KxsLeaderAmountLogService kxsLeaderAmountLogService; private final KxsLeaderAccountLogService kxsLeaderAccountLogService; private final KxsUserReRealService kxsUserReRealService; private final KxsUserPresetBeforeService kxsUserPresetBeforeService; 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 KxsMachineTransferService kxsMachineTransferService; private final KxsMachineTrackService kxsMachineTrackService; private final KxsShopOrderService kxsShopOrderService; //仓库模块 private final KxsWarehouseService kxsWarehouseService; private final KxsMachineRecycleService kxsMachineRecycleService; private final KxsMachineApplyService kxsMachineApplyService; private final KxsMachineAdvanceService kxsMachineAdvanceService; private final KxsWarehouseMachineApplyService kxsWarehouseMachineApplyService; private final KxsWarehouseStockLogService kxsWarehouseStockLogService; private final KxsWarehouseLimitLogService kxsWarehouseLimitLogService; private final KxsActTotalService kxsActTotalService; //统计模块 private final KxsUserTradeService kxsUserTradeService; private final KxsUserTradeBeforeService kxsUserTradeBeforeService; private final KxsUserTradeAfterService kxsUserTradeAfterService; private final KxsUserActTradeService kxsUserActTradeService; private final KxsMerchantTradeService kxsMerchantTradeService; private final KxsUserLogoutTradeService kxsUserLogoutTradeService; private final KxsUserGdTradeService kxsUserGdTradeService; private final KxsZlbTradeService kxsZlbTradeService; private final KxsUserNewTradeService kxsUserNewTradeService; private final KxsBannerService kxsBannerService; private final KxsColService kxsColService; @Value("${spring.profiles.active}") private String active; @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("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 ("UserRankWhiteBefore".equals(dmlData.getTableName())) { kxsUserPresetBeforeService.changeData(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()) && ("dev".equals(active) || "test".equals(active) ? "KxsMainServer":"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()) && ("dev".equals(active) || "test".equals(active) ? "KxsMainServer":"KxsProfitServer").equals(dmlData.getDatabaseName())) { kxsUserWithdrawalService.changeData(dmlData); } //盟主表 if ("Leaders".equals(dmlData.getTableName())) { kxsLeaderService.changeData(dmlData); } //盟主金额变动记录表 if ("LeaderReserveRecord".equals(dmlData.getTableName())) { kxsLeaderAmountLogService.changeData(dmlData); //盟主运营中心合并的新表,此表只同步盟主储蓄金 kxsLeaderAccountLogService.changeData(dmlData); } //盟主可提现金额变动记录表 if ("LeaderAccountRecord".equals(dmlData.getTableName())) { kxsLeaderAccountLogService.change2Data(dmlData); } //合伙人表 if ("OperatorRankWhite".equals(dmlData.getTableName())) { kxsPartnerService.changeData(dmlData); //盟主表 kxsLeaderService.changeLeaderData(dmlData); } //合伙人账户 -> 合伙人表 if ("UserAccount".equals(dmlData.getTableName()) && "KxsOpServer".equals(dmlData.getDatabaseName())) { kxsPartnerService.changeAccountData(dmlData); kxsLeaderService.changeAccountData(dmlData); } //合伙人账户余额变动记录表 if ("AmountRecordNew".equals(dmlData.getTableName()) && "KxsOpServer".equals(dmlData.getDatabaseName())) { kxsPartnerService.changeAccountLogData(dmlData); //盟主金额变动记录 kxsLeaderAccountLogService.changePartnerData(dmlData); } //用户结算卡变更记录 if ("SettlementCardChangeRecord".equals(dmlData.getTableName())) { kxsBankChangeLogService.changeData(dmlData); } //分期表数据同步 if ("ToChargeByStage".equals(dmlData.getTableName())) { kxsUserByStageService.changeData(dmlData); } //分期期数表数据同步 if ("ToChargeBackRecordSub".equals(dmlData.getTableName())) { kxsUserByStageService.changeInfoData(dmlData); } //创客重新实名开放表 if ("UserSetUnAuthRecord".equals(dmlData.getTableName())) { kxsUserReRealService.changeData(dmlData); } //商品表同步 if ("Products".equals(dmlData.getTableName())) { kxsShopGoodsService.changeData(dmlData); } //商品订单同步 if ("Orders".equals(dmlData.getTableName())) { kxsShopOrderService.changeData(dmlData); } //商品订单详情 if ("OrderProduct".equals(dmlData.getTableName())) { kxsShopOrderService.changeInfoData(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 ("PosMachinesTwo".equals(dmlData.getTableName())) { kxsMachineService.changeData(dmlData); } //未绑定划拨记录表 if ("UserStoreChange".equals(dmlData.getTableName())) { kxsMachineTransferService.changeData(dmlData); } //机具轨迹记录表 if ("StoreChangeHistory".equals(dmlData.getTableName())) { kxsMachineTrackService.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); } //仓库表 if ("StoreHouse".equals(dmlData.getTableName())) { kxsWarehouseService.changeData(dmlData); } //机具申请订单表 if ("MachineApply".equals(dmlData.getTableName())) { kxsMachineApplyService.changeData(dmlData); } //仓库预发机 if ("PreSendStockDetail".equals(dmlData.getTableName())) { kxsMachineAdvanceService.changeData(dmlData); } //分仓机具申请表 if ("StoreMachineApply".equals(dmlData.getTableName())) { kxsWarehouseMachineApplyService.changeData(dmlData); } //仓库担保记录 if ("StoreHouseAmountPromiss".equals(dmlData.getTableName())) { kxsWarehouseService.changeCreditData(dmlData); } //仓库库存变动记录表 if ("StoreBalance".equals(dmlData.getTableName())) { kxsWarehouseStockLogService.changeData(dmlData); } //仓库额度变动记录表 if ("StoreHouseAmountRecord".equals(dmlData.getTableName())) { kxsWarehouseLimitLogService.changeData(dmlData); } //小分仓额度变动记录表 if ("PreAmountRecord".equals(dmlData.getTableName())) { kxsWarehouseLimitLogService.changeSmallData(dmlData); } //仓库激活统计表 if ("StoreSnActivateSummary".equals(dmlData.getTableName())) { kxsActTotalService.changeData(dmlData); } //仓库额度表 if ("UserAccount".equals(dmlData.getTableName()) && ("dev".equals(active) || "test".equals(active) ? "KxsMainServer":"KxsProfitServer").equals(dmlData.getDatabaseName())) { kxsWarehouseService.changeAccountData(dmlData); } //小分仓库额度表 if ("UserAccount".equals(dmlData.getTableName()) && ("dev".equals(active) || "test".equals(active) ? "KxsMainServer":"KxsProfitServer").equals(dmlData.getDatabaseName())) { kxsWarehouseService.changeAccountMinData(dmlData); } //广告分类 if ("Col".equals(dmlData.getTableName()) && "KxsBsServer".equals(dmlData.getDatabaseName())) { kxsColService.changeData(dmlData); } if ("Advertisment".equals(dmlData.getTableName()) && "KxsBsServer".equals(dmlData.getDatabaseName())) { kxsBannerService.changeData(dmlData); } //创客消息表 if ("MsgPersonal".equals(dmlData.getTableName())) { kxsSysMsgService.changeUserData(dmlData); } //系统消息表 if ("MsgPlacard".equals(dmlData.getTableName())) { kxsSysMsgService.changeSysData(dmlData); } if ("MsgPlacardRead".equals(dmlData.getTableName())) { kxsSysMsgService.changeSysReadData(dmlData); } //系统弹窗表 if ("MsgAlert".equals(dmlData.getTableName())) { kxsSysMsgService.changeModalData(dmlData); } //创客交易额统计表 if ("TradeDaySummary".equals(dmlData.getTableName())) { kxsUserTradeService.changeData(dmlData); } if ("TradeDaySummary2".equals(dmlData.getTableName())) { kxsUserTradeService.change2Data(dmlData); } //存量 if ("TradeDaySummaryBefore".equals(dmlData.getTableName())) { kxsUserTradeBeforeService.changeData(dmlData); } if ("TradeDaySummary2Before".equals(dmlData.getTableName())) { kxsUserTradeBeforeService.change2Data(dmlData); } //增量 if ("TradeDaySummaryAfter".equals(dmlData.getTableName())) { kxsUserTradeAfterService.changeData(dmlData); } if ("TradeDaySummary2After".equals(dmlData.getTableName())) { kxsUserTradeAfterService.change2Data(dmlData); } //创客激活统计 if ("UserTradeMonthSummary".equals(dmlData.getTableName())) { kxsUserActTradeService.changeData(dmlData); } //商户交易额统计 if ("PosMerchantTradeSummay".equals(dmlData.getTableName())) { kxsMerchantTradeService.changeData(dmlData); } //广电注销统计表 if ("UserSimActSummary".equals(dmlData.getTableName())) { kxsUserLogoutTradeService.changeData(dmlData); } //广电话费统计表 if ("SimCardDaySummary".equals(dmlData.getTableName())) { kxsUserGdTradeService.changeData(dmlData); } //助力宝统计 if ("HelpProfitUserTradeSummay".equals(dmlData.getTableName())) { kxsZlbTradeService.changeData(dmlData); } //创客新增统计表 if ("PullnewSummary".equals(dmlData.getTableName())) { kxsUserNewTradeService.changeData(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); } } }