Spring Boot中如何按指定次数运行方法 实现固定次数发送RabbitMQ消息
实现方案
你可以选择以下任意一种方式实现指定次数发送消息后自动停止的需求:
方案1:改造现有@Scheduled逻辑,增加终止判断
你无需完全舍弃@Scheduled注解,只需要在现有逻辑中增加次数校验,发够阈值后终止任务即可:
- 先注入定时任务处理器到你的Sender类:
@Autowired private ScheduledAnnotationBeanPostProcessor scheduledAnnotationBeanPostProcessor; // 定义总发送次数,可根据实际需求修改,这里示例取预存的UUID总数 private static final int TOTAL_SEND_COUNT = uuids.size();
- 修改sendNumbers方法增加判断逻辑:
@Scheduled(initialDelay = 1000,fixedDelay = 1500) public void sendNumbers(){ // 已达发送上限直接跳过执行 if (atomicInteger.get() >= TOTAL_SEND_COUNT) { return; } int index = atomicInteger.getAndIncrement(); UUID uuid = uuids.get(index); Pair<Integer,Integer> pair = sumPair.get(uuid); MessagePostProcessor messagePostProcessor = message -> { MessageProperties messageProperties = message.getMessageProperties(); messageProperties.setCorrelationId(uuid.toString()); messageProperties.setReplyTo(response.getName()); return message; }; rabbitTemplate.convertAndSend(directExchange.getName(),routingKey,pair,messagePostProcessor); // 发送完最后一条后直接销毁该定时任务,避免后续空跑占用资源 if (index == TOTAL_SEND_COUNT - 1) { scheduledAnnotationBeanPostProcessor.postProcessBeforeDestruction(this, "sendNumbers"); } }
方案2:使用TaskScheduler手动控制调度
如果你需要更灵活的触发时机(比如收到某个请求后才开始发送),可以抛弃@Scheduled注解,手动提交调度任务:
- 注入TaskScheduler到Sender类:
@Autowired private TaskScheduler taskScheduler; private ScheduledFuture<?> sendTaskFuture;
- 编写发送任务启动方法:
/** * 启动指定次数的发送任务 * @param totalSendCount 总发送次数 */ public void startSendTask(int totalSendCount) { AtomicInteger counter = new AtomicInteger(0); // 配置初始延迟1s、间隔1.5s执行,和原有@Scheduled参数保持一致 sendTaskFuture = taskScheduler.scheduleWithFixedDelay( () -> { int currentIndex = counter.get(); if (currentIndex >= totalSendCount) { // 达到发送上限直接取消任务 sendTaskFuture.cancel(false); return; } // 原有发送逻辑 int index = counter.getAndIncrement(); UUID uuid = uuids.get(index); Pair<Integer,Integer> pair = sumPair.get(uuid); MessagePostProcessor messagePostProcessor = message -> { MessageProperties messageProperties = message.getMessageProperties(); messageProperties.setCorrelationId(uuid.toString()); messageProperties.setReplyTo(response.getName()); return message; }; rabbitTemplate.convertAndSend(directExchange.getName(),routingKey,pair,messagePostProcessor); }, new Date(System.currentTimeMillis() + 1000), 1500 ); }
方案3:简单循环实现(适合测试场景)
如果是本地测试场景,不需要异步调度的能力,直接用循环加休眠即可实现:
/** * 同步执行指定次数的消息发送 * @param totalSendCount 总发送次数 * @throws InterruptedException 线程中断异常 */ public void sendMultiTimes(int totalSendCount) throws InterruptedException { for (int i = 0; i < totalSendCount; i++) { UUID uuid = uuids.get(i); Pair<Integer,Integer> pair = sumPair.get(uuid); MessagePostProcessor messagePostProcessor = message -> { MessageProperties messageProperties = message.getMessageProperties(); messageProperties.setCorrelationId(uuid.toString()); messageProperties.setReplyTo(response.getName()); return message; }; rabbitTemplate.convertAndSend(directExchange.getName(),routingKey,pair,messagePostProcessor); // 最后一次发送不需要休眠 if (i != totalSendCount - 1) { Thread.sleep(1500); } } }
内容的提问来源于stack exchange,提问作者DozezQuest
相关产品推荐
相关产品推荐

