事件驱动微服务中消息队列宕机的处理及编程选型咨询
消息队列宕机场景下:线程编程 vs 响应式编程的选型与实现
针对你提到的「消息队列M宕机时,服务A本地存消息+定期检测M状态恢复后批量发送」的场景,我来拆解下两种编程模型的适用情况和具体实现:
一、选型分析
1. 线程编程(传统阻塞式)
适合中小流量、团队对响应式技术栈不熟悉的场景:
- 优势:逻辑直观,代码容易编写和调试,团队学习成本低,适合快速落地简单的故障恢复逻辑。
- 劣势:如果待发送的消息量极大,线程池容易出现阻塞,资源利用率偏低,难以支撑高并发的消息积压处理。
2. 响应式编程(非阻塞式)
适合高流量、高吞吐量、需要最大化资源利用率的场景:
- 优势:基于非阻塞异步模型,用少量线程就能处理大量并发任务,即使消息积压严重也能保持系统的响应性,资源利用率更高。
- 劣势:学习曲线较陡,调试和排查问题的复杂度更高,需要团队熟悉Reactor/RxJava这类响应式框架。
二、代码实现示例
1. 线程编程实现(Spring Boot + 线程池)
这里用Spring Boot的定时任务+线程池来实现检测和消息发送逻辑:
核心依赖(pom.xml)
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-jdbc</artifactId> </dependency> <dependency> <groupId>com.h2database</groupId> <artifactId>h2</artifactId> <scope>runtime</scope> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-quartz</artifactId> </dependency> </dependencies>
核心代码
import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.List; import java.util.Map; @Component public class MessageQueueRecoveryService { private final JdbcTemplate jdbcTemplate; private final MessageQueueClient mqClient; public MessageQueueRecoveryService(JdbcTemplate jdbcTemplate, MessageQueueClient mqClient) { this.jdbcTemplate = jdbcTemplate; this.mqClient = mqClient; } // 模拟服务A发送消息的逻辑:如果MQ不可用,存本地DB public void sendMessage(String message) { if (!mqClient.isAlive()) { jdbcTemplate.update("INSERT INTO pending_messages(content) VALUES (?)", message); return; } mqClient.send(message); } // 每10秒检测MQ状态,恢复后批量发送消息 @Scheduled(fixedRate = 10000) public void checkAndRecoverMessages() { if (!mqClient.isAlive()) { return; } // 批量读取本地待发送消息 List<Map<String, Object>> pendingMessages = jdbcTemplate.queryForList("SELECT id, content FROM pending_messages LIMIT 100"); if (pendingMessages.isEmpty()) { return; } // 批量发送到MQ pendingMessages.forEach(msg -> { String content = (String) msg.get("content"); mqClient.send(content); // 发送成功后删除本地记录 jdbcTemplate.update("DELETE FROM pending_messages WHERE id = ?", msg.get("id")); }); } } // 模拟MQ客户端 class MessageQueueClient { private boolean isAlive = true; // 模拟宕机时设为false public boolean isAlive() { return isAlive; } public void send(String message) { // 实际MQ发送逻辑 System.out.println("发送消息到MQ:" + message); } // 模拟宕机/恢复的方法 public void setAlive(boolean alive) { isAlive = alive; } }
2. 响应式编程实现(Spring WebFlux + R2DBC)
用Spring WebFlux的响应式模型,结合R2DBC操作数据库,实现非阻塞的检测和消息发送:
核心依赖(pom.xml)
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-webflux</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-r2dbc</artifactId> </dependency> <dependency> <groupId>io.r2dbc</groupId> <artifactId>r2dbc-h2</artifactId> <scope>runtime</scope> </dependency> </dependencies>
核心代码
import org.springframework.data.r2dbc.core.R2dbcEntityTemplate; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Component @EnableScheduling public class ReactiveMessageQueueRecoveryService { private final R2dbcEntityTemplate r2dbcEntityTemplate; private final ReactiveMessageQueueClient mqClient; public ReactiveMessageQueueRecoveryService(R2dbcEntityTemplate r2dbcEntityTemplate, ReactiveMessageQueueClient mqClient) { this.r2dbcEntityTemplate = r2dbcEntityTemplate; this.mqClient = mqClient; } // 响应式发送消息:MQ不可用则存本地 public Mono<Void> sendMessage(String message) { return mqClient.isAlive() .flatMap(alive -> { if (alive) { return mqClient.send(message); } else { return r2dbcEntityTemplate.insert(PendingMessage.class) .using(new PendingMessage(null, message)) .then(); } }); } // 定时检测MQ状态,恢复后批量发送 @Scheduled(fixedRate = 10000) public void checkAndRecoverMessages() { mqClient.isAlive() .filter(alive -> alive) .flatMapMany(alive -> r2dbcEntityTemplate.select(PendingMessage.class).limit(100)) .flatMap(msg -> mqClient.send(msg.getContent()) .then(r2dbcEntityTemplate.delete(PendingMessage.class) .matching(query -> query.where("id").is(msg.getId())) .then())) .subscribe(); } } // 待发送消息实体 class PendingMessage { private Long id; private String content; public PendingMessage(Long id, String content) { this.id = id; this.content = content; } // getter/setter public Long getId() { return id; } public void setId(Long id) { this.id = id; } public String getContent() { return content; } public void setContent(String content) { this.content = content; } } // 响应式MQ客户端 class ReactiveMessageQueueClient { private boolean isAlive = true; public Mono<Boolean> isAlive() { // 模拟异步检测MQ状态 return Mono.just(isAlive); } public Mono<Void> send(String message) { // 模拟异步发送消息到MQ return Mono.fromRunnable(() -> System.out.println("响应式发送消息到MQ:" + message)) .then(); } public void setAlive(boolean alive) { isAlive = alive; } }
总结
如果你的服务流量不大,团队更熟悉传统Java线程模型,线程编程是更稳妥的选择;如果需要处理高并发消息积压,追求更高的资源利用率,响应式编程会更合适。
内容的提问来源于stack exchange,提问作者Prash
相关产品推荐
相关产品推荐

