You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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);
        }
    }
}

问题原因分析

  1. @Async注解未生效导致主线程阻塞:UdpServer类上添加了@Async,但如果SpringBoot启动类未添加@EnableAsync注解,@Async将不生效。此时CommandLineRunner的run方法会在主线程执行,其中cf2.channel().closeFuture().sync()是同步阻塞方法,会卡住主线程,导致Spring无法完成后续的RabbitMQ监听器容器初始化,最终队列无消费者。
  2. 队列名称不匹配:RabbitConfig中创建队列时使用downlinkUdpConfig.getQueue()动态获取名称,而监听器@RabbitListener(queues = "downlink_udp_queue")是硬编码名称,若配置文件中队列名称与硬编码不一致,会导致监听器监听的队列与实际创建的队列不匹配,显示无消费者。
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 04:45:44