2026/7/30 11:01:18

大营销平台 —— 中奖记录入库业务和定时任务巡检MQ实现

大营销平台 —— 中奖记录入库业务和定时任务巡检MQ实现 一、前言回顾前面完成的领域SKU订单领域、抽奖订单领域、抽奖策略领域经过了这些领域我们已经拿到了奖品接下来就是要去实现让这个奖品入这个就需要创建一个奖品记录了。我们选择用MQ异步发放奖品为了避免MQ宕机的意外发生我们还将通过Task定时任务定时巡检将发送失败的中奖记录重发直到保证所有的任务都被发送到监听器。二、保存中奖记录的业务实现1.业务结构自然的这是一个新的业务所以我们将划分出一个新的award领域来实现领域结构如下同时刚刚我们说了Task会定时巡检MQ所以这里给定时任务也单独划分一个领域因为我们为了扩展性的考虑希望Task不止为award领域中的中奖记录进行巡检后续可能还有需要巡检的业务因此要将Task解耦出来。2.发奖业务实现这里我们没有再使用抽象类对业务骨架进行编排了因为这个业务实现起来很简单主要行为其实就只有入库没有那么多的校验操作。/** * author 印东升 * description * create 2026-07-28 11:51 */ Service public class AwardService implements IAwardService { Resource private IAwardRepository awardRepository; Resource private SendAwardMessageEvent sendAwardMessageEvent; Override public void saveUserAwardRecord(UserAwardRecordEntity userAwardRecordEntity) { //构建消息对象 SendAwardMessageEvent.SendAwardMessage sendAwardMessage new SendAwardMessageEvent.SendAwardMessage(); sendAwardMessage.setAwardTitle(userAwardRecordEntity.getAwardTitle()); sendAwardMessage.setAwardId(userAwardRecordEntity.getAwardId()); sendAwardMessage.setUserId(userAwardRecordEntity.getUserId()); BaseEvent.EventMessageSendAwardMessageEvent.SendAwardMessage sendAwardMessageEventMessage sendAwardMessageEvent.buildEventMessage(sendAwardMessage); //构建任务对象 TaskEntity taskEntity new TaskEntity(); taskEntity.setTopic(sendAwardMessageEvent.topic()); taskEntity.setState(TaskStateVO.create); taskEntity.setMessageId(sendAwardMessageEventMessage.getId()); taskEntity.setMessage(sendAwardMessageEventMessage); taskEntity.setUserId(userAwardRecordEntity.getUserId()); //构建聚合对象 UserAwardRecordAggregate userAwardRecordAggregate UserAwardRecordAggregate.builder() .taskEntity(taskEntity) .userAwardRecordEntity(userAwardRecordEntity) .build(); //保存聚合对象 awardRepository.saveUserAwardRecord(userAwardRecordAggregate); } }对于MQ的消息我们应该遵循一个全局的规范约束而这个规范在前面我们已经使用过了就是BaseEvent类在award领域中尽管具体的data和其他不一样但是为了统一格式我们需要让这个发奖消息类继承这个规范同时定义出属于发奖消息的data也就是SendAwardMessage。/** * author 印东升 * description * create 2026-07-28 11:33 */ Service public class SendAwardMessageEvent extends BaseEventSendAwardMessageEvent.SendAwardMessage { Value(${spring.rabbitmq.topic.send_award}) private String topic; Override public EventMessageSendAwardMessage buildEventMessage(SendAwardMessage data) { return EventMessage.SendAwardMessagebuilder() .id(RandomStringUtils.randomNumeric(11)) .timestamp(new Date()) .data(data) .build(); } Override public String topic() { return topic; } Data Builder NoArgsConstructor AllArgsConstructor public static class SendAwardMessage { /** * 用户ID */ private String userId; /** * 奖品ID */ private Integer awardId; /** * 奖品名称 */ private String awardTitle; } }聚合对象/** * author 印东升 * description 用户中奖记录聚合对象 * create 2026-07-28 11:48 */ Data Builder NoArgsConstructor AllArgsConstructor public class UserAwardRecordAggregate { private UserAwardRecordEntity userAwardRecordEntity; private TaskEntity taskEntity; }而真正的发消息、入库都是在仓储中实现的如果出现异常会马上被catch然后在Task表中被标记失败供后续定时任务巡检重试。/** * author 印东升 * description 奖品仓储服务 * create 2026-07-28 11:53 */ Repository Slf4j public class AwardRepository implements IAwardRepository { Resource private ITaskDao taskDao; Resource private IUserAwardRecordDao userAwardRecordDao; Resource private TransactionTemplate transactionTemplate; Resource private EventPublisher eventPublisher; Resource private IDBRouterStrategy dbRouter; Override public void saveUserAwardRecord(UserAwardRecordAggregate userAwardRecordAggregate) { UserAwardRecordEntity userAwardRecordEntity userAwardRecordAggregate.getUserAwardRecordEntity(); TaskEntity taskEntity userAwardRecordAggregate.getTaskEntity(); Integer awardId userAwardRecordEntity.getAwardId(); Long activityId userAwardRecordEntity.getActivityId(); String userId userAwardRecordEntity.getUserId(); UserAwardRecord userAwardRecord new UserAwardRecord(); userAwardRecord.setUserId(userAwardRecordEntity.getUserId()); userAwardRecord.setActivityId(userAwardRecordEntity.getActivityId()); userAwardRecord.setStrategyId(userAwardRecordEntity.getStrategyId()); userAwardRecord.setOrderId(userAwardRecordEntity.getOrderId()); userAwardRecord.setAwardId(userAwardRecordEntity.getAwardId()); userAwardRecord.setAwardTitle(userAwardRecordEntity.getAwardTitle()); userAwardRecord.setAwardTime(userAwardRecordEntity.getAwardTime()); userAwardRecord.setAwardState(userAwardRecordEntity.getAwardState().getCode()); Task task new Task(); task.setUserId(taskEntity.getUserId()); task.setTopic(taskEntity.getTopic()); task.setMessageId(taskEntity.getMessageId()); task.setMessage(JSON.toJSONString(taskEntity.getMessage())); task.setState(JSON.toJSONString(taskEntity.getState())); try { dbRouter.doRouter(userId); transactionTemplate.execute(status - { try { //写入记录 userAwardRecordDao.insert(userAwardRecord); //写入任务 taskDao.insert(task); return 1; } catch (DuplicateKeyException e) { log.error(写入中奖记录唯一索引冲突 userId:{} activityId:{} awardId:{}, userId, activityId, awardId); throw new AppException(ResponseCode.INDEX_DUP.getCode(),e); } }); } finally { dbRouter.clear(); } try{ //发送消息 eventPublisher.publish(task.getTopic(),task.getMessage()); //更新数据库记录 taskDao.updateTaskSendMessageCompleted(task); }catch (Exception e){ log.error(写入中奖记录发送MQ消息失败 userId:{} topic:{},userId,task.getTopic()); taskDao.updateTaskSendMessageFail(task); } } }如果一切正常这边的监听器消费者就应该监听到奖品发送消息了。/** * author 印东升 * description 用户奖品记录消息消费者 * create 2026-07-28 15:21 */ Slf4j Component public class SendAwardCustomer { Value(${spring.rabbitmq.topic.send_award}) private String topic; RabbitListener(queuesToDeclare Queue(value ${spring.rabbitmq.topic.send_award})) public void listener(String message){ try{ log.info(监听用户奖品发送消息 topic:{} message:{},topic,message); }catch (Exception e){ log.error(监听用户奖品发送消息消费失败 topic:{} message{},topic,message); throw e; } } }3.巡检业务实现其实就是一个通用的定时任务服务提供一些巡检方法供定时任务使用/** * author 印东升 * description 消息任务服务 * create 2026-07-28 14:37 */ Service public class TaskService implements ITaskService { Resource private ITaskRepository taskRepository; Override public ListTaskEntity queryNoSendMessageTaskList() { return taskRepository.queryNoSendMessageTaskList(); } Override public void sendMessage(TaskEntity taskEntity) { taskRepository.sendMessage(taskEntity); } Override public void updateTaskSendMessageCompleted(String userId, String messageId) { taskRepository.updateTaskSendMessageCompleted(userId,messageId); } Override public void updateTaskSendMessageFail(String userId, String messageId) { taskRepository.updateTaskSendMessageFail(userId,messageId); } }这里的主要难度在于怎么去处理分库分表其实里面的逻辑很简单就是一个重试过程成功了就把任务状态改成completed失败就改成fail然后继续重试。/** * author 印东升 * description 发送MQ消息任务队列 * create 2026-07-28 12:39 */ Slf4j Component() public class SendMessageTaskJob { Resource private ITaskService taskService; Resource private ThreadPoolExecutor executor; Resource private IDBRouterStrategy dbRouter; Scheduled(cron 0/5 * * * * ?) public void exec() { try { //获取分库数量 int dbCount dbRouter.dbCount(); for (int dbIdx 1; dbIdx dbCount; dbIdx) { int finalDbIdx dbIdx; executor.execute(() - { try { dbRouter.setDBKey(finalDbIdx); dbRouter.setTBKey(0); ListTaskEntity taskEntities taskService.queryNoSendMessageTaskList(); if (taskEntities.isEmpty()) return; //发送MQ消息 for (TaskEntity taskEntity : taskEntities) { //开启线程发送提高发送效率。 executor.execute(() - { try { taskService.sendMessage(taskEntity); taskService.updateTaskSendMessageCompleted(taskEntity.getUserId(), taskEntity.getMessageId()); } catch (Exception e) { log.info(定时任务发送MQ消息失败 userId:{} topic:{} , taskEntity.getUserId(), taskEntity.getTopic()); taskService.updateTaskSendMessageFail(taskEntity.getUserId(), taskEntity.getMessageId()); } }); } } finally { dbRouter.clear(); } }); } } catch (Exception e) { log.error(定时任务扫描MQ任务表发送消息失败); } } }三、测试/** * author Fuzhengwei bugstack.cn 小傅哥 * description 奖品服务测试 * create 2024-04-06 11:27 */ Slf4j RunWith(SpringRunner.class) SpringBootTest public class AwardServiceTest { Resource private IAwardService awardService; /** * 模拟发放抽奖记录流程中会发送MQ以及接收MQ消息还有 task 表补偿发送MQ */ Test public void test_saveUserAwardRecord() throws InterruptedException { for (int i 0; i 100; i) { UserAwardRecordEntity userAwardRecordEntity new UserAwardRecordEntity(); userAwardRecordEntity.setUserId(xiaofuge); userAwardRecordEntity.setActivityId(100301L); userAwardRecordEntity.setStrategyId(100006L); userAwardRecordEntity.setOrderId(RandomStringUtils.randomNumeric(12)); userAwardRecordEntity.setAwardId(101); userAwardRecordEntity.setAwardTitle(OpenAI 增加使用次数); userAwardRecordEntity.setAwardTime(new Date()); userAwardRecordEntity.setAwardState(AwardStateVO.create); awardService.saveUserAwardRecord(userAwardRecordEntity); Thread.sleep(500); } new CountDownLatch(1).await(); } }100条数据发完后就不发了