概述
在现代微服务架构中,消息队列作为重要的组件被广泛应用于解耦系统间的数据传输。spring boot提供了强大的支持来集成各种消息中间件,如rabbitmq、kafka、rocketmq等。本文将详细介绍如何在spring boot启动类中配置消费端,使其能够随着服务启动自动开始循环消费消息。
1. 基础环境准备
1.1 maven依赖配置
首先需要在项目的 [pom.xml](file://d:\lshm\draco-center\draco-api\pom.xml) 文件中添加相应的依赖:xml
<dependencies> <!-- spring boot starter --> <dependency> <groupid>org.springframework.boot</groupid> <artifactid>spring-boot-starter</artifactid> </dependency>
<!-- spring boot web starter (可选) -->
<dependency>
<groupid>org.springframework.boot</groupid>
<artifactid>spring-boot-starter-web</artifactid>
</dependency>
<!-- spring boot amqp (rabbitmq) -->
<dependency>
<groupid>org.springframework.boot</groupid>
<artifactid>spring-boot-starter-amqp</artifactid>
</dependency>
<!-- 或者 kafka 依赖 -->
<dependency>
<groupid>org.springframework.kafka</groupid>
<artifactid>spring-kafka</artifactid>
</dependency>
</dependencies>
1.2 配置文件设置
在 application.yml 中配置消息队列连接参数:yaml
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest virtual-host: / kafka: bootstrap-servers: localhost:9092 consumer: group-id: my-consumer-group
2. 消息消费者实现
2.1 rabbitmq消费者实现
创建一个基本的消息消费者类:
@component @slf4j public class messageconsumer {
@rabbitlistener(queues = "my.queue.name")
public void handlemessage(string message) {
log.info("接收到消息: {}", message);
// 处理业务逻辑
processmessage(message);
}
private void processmessage(string message) {
try {
// 模拟业务处理
thread.sleep(1000);
log.info("消息处理完成: {}", message);
} catch (interruptedexception e) {
log.error("处理消息时发生异常", e);
thread.currentthread().interrupt();
}
}
}
2.2 kafka消费者实现
对于kafka消费者,可以这样实现:
@component @slf4j public class kafkamessageconsumer {
@kafkalistener(topics = "my-topic", groupid = "my-consumer-group")
public void listen(string message) {
log.info("kafka接收到消息: {}", message);
processmessage(message);
}
private void processmessage(string message) {
try {
// 模拟业务处理
log.info("处理kafka消息: {}", message);
} catch (exception e) {
log.error("处理kafka消息时发生异常", e);
}
}
}
3. 启动类配置
3.1 基本启动类结构
创建spring boot启动类,并配置消息消费者的自动启动:
@springbootapplication @enablerabbit public class messageconsumerapplication {
public static void main(string[] args) {
springapplication app = new springapplication(messageconsumerapplication.class);
app.setwebapplicationtype(webapplicationtype.none); // 如果不需要web功能
app.run(args);
}
}
3.2 自定义启动监听器
为了更好地控制消费者的启动时机,可以实现 applicationrunner 接口:
@component @slf4j public class consumerstartuprunner implements applicationrunner {
@autowired
private messageconsumerservice messageconsumerservice;
@override
public void run(applicationarguments args) throws exception {
log.info("应用程序启动完成,开始初始化消息消费者");
messageconsumerservice.startconsuming();
}
}
4. 消费者服务管理
4.1 消费者服务接口设计
创建一个消费者服务接口来统一管理消费者的生命周期:
public interface messageconsumerservice { /** * 开始消费消息 */ void startconsuming();
/**
* 停止消费消息
*/
void stopconsuming();
/**
* 获取消费者状态
* @return 消费者状态
*/
consumerstatus getstatus();
}
4.2 消费者状态枚举
定义消费者状态枚举:
public enum consumerstatus { /** * 初始化状态 */ initialized,
/**
* 正在运行
*/
running,
/**
* 已停止
*/
stopped,
/**
* 错误状态
*/
error
}
4.3 消费者服务实现
实现消费者服务的具体逻辑:
@service @slf4j public class messageconsumerserviceimpl implements messageconsumerservice {
private volatile consumerstatus status = consumerstatus.initialized;
private final executorservice executorservice = executors.newfixedthreadpool(5);
private volatile boolean running = false;
@autowired
private rabbittemplate rabbittemplate;
@override
public void startconsuming() {
if (status == consumerstatus.running) {
log.warn("消费者已在运行中");
return;
}
running = true;
status = consumerstatus.running;
// 启动多个消费者线程
for (int i = 0; i < 3; i++) {
executorservice.submit(new messageconsumertask(i));
}
log.info("消息消费者已启动");
}
@override
public void stopconsuming() {
running = false;
status = consumerstatus.stopped;
executorservice.shutdown();
try {
if (!executorservice.awaittermination(60, timeunit.seconds)) {
executorservice.shutdownnow();
}
} catch (interruptedexception e) {
executorservice.shutdownnow();
thread.currentthread().interrupt();
}
log.info("消息消费者已停止");
}
@override
public consumerstatus getstatus() {
return status;
}
/**
* 消息消费任务类
*/
private class messageconsumertask implements runnable {
private final int taskid;
public messageconsumertask(int taskid) {
this.taskid = taskid;
}
@override
public void run() {
log.info("消费者任务 {} 已启动", taskid);
while (running && !thread.currentthread().isinterrupted()) {
try {
// 这里模拟从队列获取消息并处理
consumemessage();
thread.sleep(1000); // 避免过度循环
} catch (interruptedexception e) {
log.info("消费者任务 {} 被中断", taskid);
thread.currentthread().interrupt();
break;
} catch (exception e) {
log.error("消费者任务 {} 处理消息时发生异常", taskid, e);
}
}
log.info("消费者任务 {} 已结束", taskid);
}
private void consumemessage() {
// 实际的消息消费逻辑
log.debug("消费者任务 {} 正在检查新消息", taskid);
// 这里应该调用具体的消费逻辑
}
}
}
5. 应用程序生命周期管理
5.1 优雅关闭配置
为了让消费者能够优雅地关闭,需要实现 disposablebean 接口或使用 @predestroy 注解:
@service @slf4j public class gracefulshutdownservice implements disposablebean {
@autowired
private messageconsumerservice messageconsumerservice;
@override
public void destroy() throws exception {
log.info("应用程序正在关闭,停止消息消费者");
messageconsumerservice.stopconsuming();
log.info("消息消费者已安全关闭");
}
}
5.2 jvm关闭钩子
也可以注册jvm关闭钩子来确保资源正确释放:
@component @slf4j public class shutdownhookregistrar implements applicationlistener<contextrefreshedevent> {
@autowired
private messageconsumerservice messageconsumerservice;
@override
public void onapplicationevent(contextrefreshedevent event) {
runtime.getruntime().addshutdownhook(new thread(() -> {
log.info("jvm关闭钩子触发,停止消息消费者");
messageconsumerservice.stopconsuming();
}));
}
}
6. 异常处理与重试机制
6.1 消息处理异常处理
为消息处理增加完善的异常处理机制:
@component @slf4j public class robustmessageconsumer {
@rabbitlistener(queues = "my.queue.name")
public void handlemessage(string message) {
int retrycount = 0;
final int maxretries = 3;
while (retrycount <= maxretries) {
try {
log.info("处理消息: {}", message);
processmessage(message);
return; // 成功处理后返回
} catch (exception e) {
retrycount++;
log.error("处理消息失败,第{}次重试", retrycount, e);
if (retrycount > maxretries) {
log.error("消息处理失败超过最大重试次数,发送到死信队列: {}", message);
sendtodeadletterqueue(message);
break;
}
// 指数退避策略
try {
long delay = (long) math.pow(2, retrycount) * 1000;
thread.sleep(delay);
} catch (interruptedexception ie) {
thread.currentthread().interrupt();
break;
}
}
}
}
private void processmessage(string message) throws exception {
// 实际的业务处理逻辑
if (message.contains("error")) {
throw new runtimeexception("模拟处理错误");
}
log.info("成功处理消息: {}", message);
}
private void sendtodeadletterqueue(string message) {
// 发送到死信队列的逻辑
log.info("发送消息到死信队列: {}", message);
}
}
6.2 监控与健康检查
添加消费者健康检查功能:
@component public class consumerhealthindicator implements healthindicator {
@autowired
private messageconsumerservice messageconsumerservice;
@override
public health health() {
consumerstatus status = messageconsumerservice.getstatus();
if (status == consumerstatus.running) {
return health.up()
.withdetail("consumerstatus", status)
.withdetail("message", "消息消费者正常运行")
.build();
} else if (status == consumerstatus.error) {
return health.down()
.withdetail("consumerstatus", status)
.withdetail("message", "消息消费者出现错误")
.build();
} else {
return health.unknown()
.withdetail("consumerstatus", status)
.withdetail("message", "消息消费者状态未知")
.build();
}
}
}
7. 高级配置选项
7.1 并发消费者配置
配置并发消费者数量:
@configuration @enablerabbit public class rabbitmqconfig {
@bean
public simplerabbitlistenercontainerfactory rabbitlistenercontainerfactory(
connectionfactory connectionfactory) {
simplerabbitlistenercontainerfactory factory =
new simplerabbitlistenercontainerfactory();
factory.setconnectionfactory(connectionfactory);
factory.setconcurrentconsumers(3); // 最小并发消费者数
factory.setmaxconcurrentconsumers(10); // 最大并发消费者数
factory.setprefetchcount(1); // 每个消费者预取的消息数
return factory;
}
}
7.2 批量消费配置
启用批量消费模式:
@component @slf4j public class batchmessageconsumer {
@rabbitlistener(queues = "batch.queue", containerfactory = "batchrabbitlistenercontainerfactory")
public void handlebatchmessages(list<string> messages) {
log.info("批量接收 {} 条消息", messages.size());
for (string message : messages) {
processmessage(message);
}
}
private void processmessage(string message) {
log.info("处理消息: {}", message);
}
}
@configuration @enablerabbit class batchrabbitmqconfig {
@bean
public simplerabbitlistenercontainerfactory batchrabbitlistenercontainerfactory(
connectionfactory connectionfactory) {
simplerabbitlistenercontainerfactory factory =
new simplerabbitlistenercontainerfactory();
factory.setconnectionfactory(connectionfactory);
factory.setbatchlistener(true); // 启用批处理
factory.setbatchsize(10); // 批处理大小
factory.setreceivetimeout(5000l); // 接收超时时间
return factory;
}
}
8. 测试验证
8.1 单元测试
编写消费者服务的单元测试:
@springboottest @testpropertysource(properties = { "spring.rabbitmq.host=localhost", "spring.rabbitmq.port=5672" }) class messageconsumerservicetest {
@autowired
private messageconsumerservice messageconsumerservice;
@test
void teststartconsuming() {
// 启动消费者
messageconsumerservice.startconsuming();
// 验证状态
assertequals(consumerstatus.running, messageconsumerservice.getstatus());
// 停止消费者
messageconsumerservice.stopconsuming();
assertequals(consumerstatus.stopped, messageconsumerservice.getstatus());
}
}
8.2 集成测试
编写完整的集成测试:
@springboottest @testpropertysource(properties = { "spring.rabbitmq.host=localhost", "spring.rabbitmq.port=5672" }) class messageintegrationtest {
@autowired
private rabbittemplate rabbittemplate;
@spybean
private messageconsumer messageconsumer;
@test
void testmessageconsumption() throws interruptedexception {
// 发送测试消息
string testmessage = "hello, world!";
rabbittemplate.convertandsend("my.queue.name", testmessage);
// 等待消息被消费
thread.sleep(2000);
// 验证消息已被消费
verify(messageconsumer, times(1)).handlemessage(testmessage);
}
}
9. 生产环境最佳实践
9.1 日志配置
配置详细的日志记录:
logging: level: com.yourpackage.consumer: debug org.springframework.amqp: info pattern: console: "%d{yyyy-mm-dd hh:mm:ss} [%thread] %-5level %logger{36} - %msg%n"
9.2 性能监控
添加性能监控指标:
@component @slf4j public class performancemonitor {
private final meterregistry meterregistry;
private final counter messagecounter;
private final timer processingtimer;
public performancemonitor(meterregistry meterregistry) {
this.meterregistry = meterregistry;
this.messagecounter = counter.builder("messages.consumed")
.description("已消费的消息总数")
.register(meterregistry);
this.processingtimer = timer.builder("message.processing.time")
.description("消息处理耗时")
.register(meterregistry);
}
public void recordmessageprocessed(long processingtimems) {
messagecounter.increment();
processingtimer.record(processingtimems, timeunit.milliseconds);
}
}
10. 总结
通过以上详细的配置和实现,我们可以在spring boot应用中成功配置消费端随服务启动循环消费消息的功能。关键要点包括:
- 正确的依赖配置:确保引入了合适的消息中间件依赖
- 合理的消费者设计:采用多线程、异常处理、重试机制等
- 完善的生命周期管理:实现优雅启动和关闭
- 健壮的异常处理:确保系统稳定性和可靠性
- 充分的测试覆盖:保证功能正确性和稳定性
这种配置方式使得消息消费者能够在应用启动时自动开始工作,并且具备良好的容错能力和监控能力,在生产环境中具有很高的实用价值。
到此这篇关于spring boot应用中配置消费端随服务启动循环消费消息的文章就介绍到这了,更多相关springboot配置消费端循环消费消息内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论