rabbitmq(广播方式)

通过休息队列控制用户上线下线

  1. 创建交换机
@Bean(name="oninechange")
	public FanoutExchange onlineExchange() {
		log.info("【交换机实例{}创建成功】", FANOUT_EXCHANGE_NAME);
		return new FanoutExchange(FANOUT_EXCHANGE_NAME);
	}

2.创建队列

@Bean(name="onlinequeue")
	public Queue onlineQueue() {
		log.info("【队列{}实例创建成功】", MYSESSIONQUEUENAME);
		return new Queue(MYSESSIONQUEUENAME, false);
	}

.

3.绑定队列到交换机

@Bean(name="binginline")
	public Binding binglineQueue1ToExchange() {
		log.info("【绑定队列{}到交换机{}成功】", MYSESSIONQUEUENAME, FANOUT_EXCHANGE_NAME);
		return BindingBuilder.bind(onlineQueue()).to(onlineExchange());
	}

4.添加消息处理方式,监听。

@Bean(name="onlinemessageContainer")
	public SimpleMessageListenerContainer messageContainer(SessionSynchronizeReceiver SessionReceiver,
			ConnectionFactory connectionFactory) {
		// 加载处理消息A的队列
		SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
		// 设置接收多个队列里面的消息,这里设置接收队列A
		// 假如想一个消费者处理多个队列里面的信息可以如下设置:
		// container.setQueues(queueA(),queueB(),queueC());
		container.setQueues(onlineQueue());
		container.setExposeListenerChannel(true);
		// 设置最大的并发的消费者数量
		container.setMaxConcurrentConsumers(10);
		// 最小的并发消费者的数量
		container.setConcurrentConsumers(1);
		// 设置确认模式手工确认
		container.setAcknowledgeMode(AcknowledgeMode.MANUAL);
		container.setMessageListener(SessionReceiver);
		return container;
	}

5.创建监听,接受消息队列信息

@Component
	public class SessionSynchronizeReceiver implements ChannelAwareMessageListener {
		@Override
		public void onMessage(Message message, com.rabbitmq.client.Channel channel) throws Exception {
			log.debug("receiver {},channel :{}", message, channel);
			try {
				String str = new String(message.getBody(), message.getMessageProperties().getContentEncoding());
				SimpleEntry<String, Object> sim = protostuffSerializeUtil.deserializeByStr(str);
				log.debug("sim  k:{} v :{}", sim.getEntity1(), sim.getEntity2());
				if (StringUtils.equals(MYSESSIONQUEUENAME, sim.getEntity1())) {
					String token =(String) sim.getValue();
					if(!StringUtils.isBlank(token)){
					//Login login = loginMapper.selectByPrimaryKey(token);
						Object obj = jcClient.get(RedisConstant.user_online_status+"@"+token);
						if(obj==null){
							final Long time = System.currentTimeMillis();
							jcClient.put(RedisConstant.user_online_status+"@"+token,time);
							loginMapper.online(token);
							return ;
						}
					}	
					return;
				}
				//log.info("synchronize session {}", Objs.toString(sim.getEntity2()));
				//sessionEngine.local().remove(Objs.toString(sim.getEntity2()));
			} catch (Exception e) {
				log.error("", e);
			} finally {
				channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
			}
		}

	}

6。创建发送消息到指定交换机(所属队列)

/** 上线 */
	public boolean online(String access_token) {
		//QueryHelper.update(" UPDATE LOGIN SET uonline ='1' WHERE access_token =?   ", access_token);
		rabbitTemplate.convertAndSend(FANOUT_EXCHANGE_NAME, "",
				protostuffSerializeUtil.serializeToStr(new SimpleEntry<>(MYSESSIONQUEUENAME, access_token)));
		return true;
	}
### RabbitMQ 实现消息广播机制及配置方法 #### 广播模式简介 为了实现消息广播,在RabbitMQ中可以采用`fanout`类型的交换机来完成这一功能。这种交换机会将接收到的消息无差别地转发给所有绑定到该交换机上的队列,而不考虑路由键。 #### 配置Fanout Exchange 创建一个名为 `broadcast_exchange` 的 fanout 类型 exchange: ```python import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明一个fanout类型的exchange channel.exchange_declare(exchange='broadcast_exchange', exchange_type='fanout') ``` #### 绑定多个消费者队列至Exchange 为了让不同消费者都能接收相同的消息副本,需分别为每个消费者声明独立的临时队列并将其与上述定义好的 `broadcast_exchange` 进行绑定: ```python result = channel.queue_declare(queue='', exclusive=True) queue_name = result.method.queue # 将新创建的队列绑定到'broadcast_exchange' channel.queue_bind(exchange='broadcast_exchange', queue=queue_name) def callback(ch, method, properties, body): print(f"Received {body}") # 开始消费来自指定队列的消息 channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) print(' [*] Waiting for messages.') channel.start_consuming() ``` 当有新的生产者发送消息时,只需向这个特定名称的exchange发布即可让所有订阅它的队列都获得一份拷贝[^1]: ```python message = "Broadcast message" channel.basic_publish(exchange='broadcast_exchange', routing_key='', body=message) ``` 通过这种方式设置后,任何连接至此系统的客户端只要绑定了相应的队列就能监听到来自于同一个源头发出的信息流,从而实现了高效可靠的消息广播服务[^2].
评论
添加红包

请填写红包祝福语或标题

红包个数最小为10个

红包金额最低5元

当前余额3.43前往充值 >
需支付:10.00
成就一亿技术人!
领取后你会自动成为博主和红包主的粉丝 规则
hope_wisdom
发出的红包
实付
使用余额支付
点击重新获取
扫码支付
钱包余额 0

抵扣说明:

1.余额是钱包充值的虚拟货币,按照1:1的比例进行支付金额的抵扣。
2.余额无法直接购买下载,可以购买VIP、付费专栏及课程。

余额充值