SpringBoot整合Netty与RabbitMQ时@RabbitListener失效问题排查
SpringBoot整合Netty与RabbitMQ时@RabbitListener失效问题排查
问题描述
在SpringBoot项目中同时整合Netty与RabbitMQ时,出现@RabbitListener注解无法生效的情况:RabbitMQ管理界面显示downlink_udp_queue队列无消费者,业务无法接收队列消息。初步怀疑与UdpServerHandler继承SimpleChannelInboundHandler有关,需排查原因并解决。
关键Maven依赖
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.2.1.RELEASE</version> <relativePath/> <!-- lookup parent from repository --> </parent> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> <java.version>1.8</java.version> <spring-cloud.version>Hoxton.RELEASE</spring-cloud.version> </properties> <dependencies> <dependency> <groupId>io.netty</groupId> <artifactId>netty-all</artifactId> <version>4.1.43.Final</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> </dependencies>
核心代码片段
Netty服务端代码
@Component @Order(value = 0) @Async public class UdpServer implements CommandLineRunner { private Logger logger = LoggerFactory.getLogger(UdpServer.class); @Value("${udp.port:8896}") private int udpPort; @Autowired UdpServerHandler udpServerHandler; private EventLoopGroup eventLoopGroup; @Override public void run(String... args) { eventLoopGroup = new NioEventLoopGroup(); try { Bootstrap bootstrap = new Bootstrap(); bootstrap.group(eventLoopGroup) .channel(NioDatagramChannel.class) .option(ChannelOption.SO_BROADCAST, true) .option(ChannelOption.SO_REUSEADDR,true) .handler(udpServerHandler); logger.info("udp server start udp port: ",udpPort); ChannelFuture cf2 = bootstrap.bind(udpPort).sync(); cf2.channel().closeFuture().sync(); } catch (Exception e) { e.printStackTrace(); } } @PreDestroy public void closeUdp(){ logger.info("udp server closed"); eventLoopGroup.shutdownGracefully(); } }
UdpServerHandler代码
@Component @SuppressWarnings("all") public class UdpServerHandler extends SimpleChannelInboundHandler<DatagramPacket> { private Logger logger = LoggerFactory.getLogger(UdpServerHandler.class); @Autowired UdpService udpService; @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { try { Channel channel = ctx.channel(); SocketAddress socketAddress = channel.remoteAddress(); if (socketAddress != null){ String channelKey = socketAddress.toString(); logger.info("channel :{} connected", channelKey); } }catch (Exception e){ e.printStackTrace(); } super.channelActive(ctx); } @Override protected void channelRead0(ChannelHandlerContext ctx, DatagramPacket datagramPacket) { try { String sender = datagramPacket.sender().toString(); logger.info(sender); }catch (Exception e){ e.printStackTrace(); } try { udpService.solveData(ctx, datagramPacket); }catch (Exception e){ e.printStackTrace(); } } @Override public void channelInactive(ChannelHandlerContext ctx) throws Exception { String channelKey = ctx.channel().remoteAddress().toString(); logger.info("channel :{} closed", channelKey); super.channelInactive(ctx); } }
RabbitMQ配置代码
@Configuration public class RabbitConfig { @Value("${spring.rabbitmq.host}") private String host; @Value("${spring.rabbitmq.port}") private int port; @Value("${spring.rabbitmq.username}") private String username; @Value("${spring.rabbitmq.password}") private String password; @Value("${spring.rabbitmq.virtual-host}") private String virtualHost; @Autowired DownlinkUdpConfig downlinkUdpConfig; @Bean public ConnectionFactory connectionFactory() { CachingConnectionFactory connectionFactory = new CachingConnectionFactory(host, port); connectionFactory.setUsername(username); connectionFactory.setPassword(password); connectionFactory.setVirtualHost(virtualHost); connectionFactory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED); return connectionFactory; } @Bean @Primary @Scope("prototype") Encoder multipartFormEncoder() { return new SpringFormEncoder(); } @Bean @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) public RabbitTemplate rabbitTemplate() { return new RabbitTemplate(connectionFactory()); } @Bean public Queue DownlinkQueue() { Map<String, Object> arguments=new HashMap<>(3); arguments.put("x-dead-letter-exchange", downlinkUdpConfig.getDeadExchange()); arguments.put("x-dead-letter-routing-key", downlinkUdpConfig.getDeadRoutingKey()); arguments.put("x-message-ttl",downlinkUdpConfig.getxMessageTtl()); return new Queue(downlinkUdpConfig.getQueue(), true,false,false,arguments); } @Bean public DirectExchange DownlinkExchange() { return new DirectExchange(downlinkUdpConfig.getExchange()); } @Bean public Binding bindingDownlinkQueue() { return BindingBuilder.bind(DownlinkQueue()).to(DownlinkExchange()).with(downlinkUdpConfig.getRoutingKey()); } @Bean public Queue DownlinkDeadQueue(){ return new Queue(downlinkUdpConfig.getDeadQueue(),true); } @Bean public DirectExchange DownlinkDeadExchange(){ return new DirectExchange(downlinkUdpConfig.getDeadExchange(),true,false); } @Bean public Binding bindingDownlinkDead(){ return BindingBuilder.bind(DownlinkDeadQueue()).to(DownlinkDeadExchange()).with(downlinkUdpConfig.getDeadRoutingKey()); } }
RabbitMQ监听器代码
@Component public class DownlinkUdpListener { Logger logger = LoggerFactory.getLogger(this.getClass()); @RabbitListener(queues = "downlink_udp_queue") public void process(Channel channel, Message message) { //1.get message from rabbitMq String rabbitMessage = new String(message.getBody()); logger.info("get result from rabbitMq is:{}", rabbitMessage); //1.1 ack get message try { channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); } catch (IOException e) { logger.info("ack error, message is:{}", rabbitMessage); } } }
问题原因分析
- @Async注解未生效导致主线程阻塞:
UdpServer类上添加了@Async,但如果SpringBoot启动类未添加@EnableAsync注解,@Async将不生效。此时CommandLineRunner的run方法会在主线程执行,其中cf2.channel().closeFuture().sync()是同步阻塞方法,会卡住主线程,导致Spring无法完成后续的RabbitMQ监听器容器初始化,最终队列无消费者。 - 队列名称不匹配:
RabbitConfig中创建队列时使用downlinkUdpConfig.getQueue()动态获取名称,而监听器@RabbitListener(queues = "downlink_udp_queue")是硬编码名称,若配置文件中队列名称与硬编码不一致,会导致监听器监听的队列与实际创建的队列不匹配,显示无消费者。 - Netty Handler的Spring管理问题:虽然
UdpServerHandler被@Component修饰,但Netty的EventLoop线程并非Spring管理线程,不过这只会影响Handler中Spring Bean的注入时机(实际上下文初始化完成后注入已完成),不会直接导致@RabbitListener失效。
解决办法
方法一:确保@Async生效
在SpringBoot启动类上添加@EnableAsync注解,让UdpServer的run方法在异步线程执行,避免阻塞主线程:
@SpringBootApplication @EnableAsync public class YourApplication { public static void main(String[] args) { SpringApplication.run(YourApplication.class, args); } }
方法二:修改Netty启动逻辑,避免阻塞主线程
去掉@Async注解,将run方法中的阻塞逻辑改为异步监听,不调用closeFuture().sync(),而是添加监听器处理关闭事件:
@Override public void run(String... args) { eventLoopGroup = new NioEventLoopGroup(); try { Bootstrap bootstrap = new Bootstrap(); bootstrap.group(eventLoopGroup) .channel(NioDatagramChannel.class) .option(ChannelOption.SO_BROADCAST, true) .option(ChannelOption.SO_REUSEADDR,true) .handler(udpServerHandler); logger.info("udp server start udp port: {}", udpPort); // 修复日志占位符问题 ChannelFuture cf2 = bootstrap.bind(udpPort).sync(); // 用监听器替代sync(),避免阻塞主线程 cf2.channel().closeFuture().addListener(future -> { eventLoopGroup.shutdownGracefully(); logger.info("udp server closed"); }); } catch (Exception e) { e.printStackTrace(); eventLoopGroup.shutdownGracefully(); } }
同时可以去掉@PreDestroy方法,因为关闭逻辑已移到监听器中。
方法三:统一队列名称
确保配置文件中downlinkUdpConfig对应的队列名称与监听器硬编码名称一致,或者监听器也使用动态配置:
// 修改监听器,使用配置注入队列名称 @Component public class DownlinkUdpListener { Logger logger = LoggerFactory.getLogger(this.getClass()); @Value("${downlink.udp.queue}") // 对应配置文件中的队列名称 private String downlinkQueue; @RabbitListener(queues = "#{downlinkQueue}") public void process(Channel channel, Message message) { // 业务逻辑不变 } }
额外优化
修复UdpServer中的日志占位符问题,将logger.info("udp server start udp port: ",udpPort);改为logger.info("udp server start udp port: {}",udpPort);,确保端口号正确输出。
内容的提问来源于stack exchange,提问作者jialing xv
相关产品推荐
相关产品推荐

