当前位置: 代码网 > it编程>编程语言>Java > RabbitMQ广播模式(动态生成queue)使用及说明

RabbitMQ广播模式(动态生成queue)使用及说明

2026年09月28日 • Java •我要评论
rabbitmq的广播机制和activemq有所不同。先来梳理下rabbitmq中消息从产生到消费的流程吧:而exchange 存在多种类型,这里就只说广播模式(fanout)了。在广播模式中,一个e

rabbitmq的广播机制和activemq有所不同。

先来梳理下rabbitmq中消息从产生到消费的流程吧:

而exchange 存在多种类型,这里就只说广播模式(fanout)了。

在广播模式中,一个exchange对应多个queue,会向每个queue都发送信息,然后不同的queue再由其对应的消费者消费信息,即完成了广播。

因为广播模式中不关注routingkey和queue,只需要queue的queue name唯一即可,所以这里把routingkey移出来了,实际上还是会经过的哦。

1.新建一个spring boot 项目并引入官方的amqp包

<dependency>
	<groupid>org.springframework.amqp</groupid>
	<artifactid>spring-rabbit</artifactid>
</dependency>

2.添加rabbitmq连接参数

spring.rabbitmq.host=xx
spring.rabbitmq.port=5672
spring.rabbitmq.username=xx
spring.rabbitmq.password=xx

3.创建生产者

同时添加一个配置参数rabbit.exchange,用来动态指定exchange(比如使用apollo或者spring cloud config)

增加配置参数

rabbit.exchange=testexchange

增加生产者

@component
public class producer {
    @autowired
    private rabbittemplate rabbittemplate;
    @value("${rabbit.exchange}")
    private string exchange;

    public void sendinfo(){
        message message = new message("123".getbytes(),new messageproperties());
        rabbittemplate.send(exchange,"",message);
    }
}

4.动态创建queue和消费者

因为需要执行createqueue方法才能生成一个queue和消费者,所以这里先用@component指定扫描当前类,再用@postconstruct指定扫描时执行该方法。

因为创建的queue是临时queue,当消费者消失时,该queue就会自动删除,因为创建queue也是由rabbitmq自行生成的,所以queue name一定是唯一的。

这样在集群部署时,就可以做到即开即用了,就算关闭了服务,对应的queue也会自动消失。

@component
public class consumer {
    @autowired
    private rabbittemplate rabbittemplate;
    @value("${rabbit.exchange}")
    private string exchange;

    @postconstruct
    public void createqueue(){
        channel channel = rabbittemplate.getconnectionfactory().createconnection().createchannel(true);

        try{
            /**
             * 与生产者使用同一个交换机
             */
            channel.exchangedeclare(exchange, "fanout",true);
            /**
             * 获取一个随机的队列名称,使用默认方式,产生的队列为临时队列,在没有消费者时将会自动删除
             */
            string queuename = channel.queuedeclare().getqueue();

            /**
             * 关联 exchange 和 queue ,因为是广播无需指定routekey,routingkey设置为空字符串
             */
            // channel.queuebind(queue, exchange, routingkey)
            channel.queuebind(queuename, exchange, "");

            com.rabbitmq.client.consumer consumer = new defaultconsumer(channel) {
                @override
                public void handledelivery(string consumertag, envelope envelope,
                                           amqp.basicproperties properties, byte[] body) throws ioexception {
                    string message = new string(body, "utf-8");

                    if(stringutils.isempty(message)){
                        return;
                    }

                    /**
                     * 对信息做操作
                     */
                }

            };
            //true 自动回复ack
            channel.basicconsume(queuename, true, consumer);
        }catch (exception ex){
        }
    }
}

5.手动在rabbitmq控制界面创建exchange

如下图选择广播模式,并且为永久exchange,如果要设置临时exchange的话要修改4中的如下语句最后一个参数为false

channel.exchangedeclare(exchange, "fanout",true);

6.启动项目,测试

在consumer中打个断点,就能看到在启动时调用createqueue后所产生的队列名

到rabbitmq控制台上搜索下,确实存在

点开看下对应的exchange

完全ok,发送条信息试下,也能够接收到

可以试下再重复第4步,再建一个queue,或者再起一个项目,你会发现两个消费者都能接收到信息。

然后我们把项目关闭,看看queue还在不在,明显已经木有了。

总结

以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。

赞 (0)

相关文章:

版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。

发表评论

验证码:
Copyright © 2017-2026  代码网 保留所有权利. 粤ICP备2024248653号
站长QQ:2386932994 | 联系邮箱:2386932994@qq.com