package com.yixin.routing;
import com.rabbitmq.client.*;
import com.yixin.util.ConnectionUtil;
import java.io.IOException;
public class Consumer2 {
//需要绑定的交换机
private final static String EXCHANGE_NAME = "direct-exchange";
//队列的名称
private final static String QUEUE_NAME = "direct-queue-2";
public static void main(String[] args) {
Connection connection = null;
Channel channel = null;
try {
// 1: 从连接工厂中获取连接
connection = connection = ConnectionUtil.getConnection("消费者2","服务器地址",5672,"/","yixin","123456");
// 2: 从连接中获取通道channel
channel = connection.createChannel();
// 3: 声明队列,如果queue已经被创建过一次了,则可以不需要定义
channel.queueDeclare(QUEUE_NAME, false, false, false, null);
//绑定队列到交换机
//4:绑定队列到交换机,指定路由key为update、delete、add(可以多个也可以单个)
channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,"select");
//5:同一时刻服务器只会发送一条消息给消费者
channel.basicQos(1);
// 6: 定义接受消息的回调
Channel finalChannel = channel;
//参数2:手动返回完成状态
finalChannel.basicConsume(QUEUE_NAME, false, new DeliverCallback() {
public void handle(String s, Delivery delivery) throws IOException {
System.out.println(QUEUE_NAME + ":外汇跟单gendan5.com收到消息是:" + new String(delivery.getBody(), "UTF-8"));
//返回确定状态,由于我们是设置手动返回状态的,如果没有调用该方法,那么消息仍然在我们队列中,不会被删除
finalChannel.basicAck(delivery.getEnvelope().getDeliveryTag(),false);
}
}, new CancelCallback() {
public void handle(String consumerTag) throws IOException {
}
});
System.out.println(QUEUE_NAME + ":开始接受消息");
//使我们的程序一直处于运行状态,也就是一直处于监听的状态,这样就不需要使用while(true)了。
System.in.read();
} catch (Exception ex) {
ex.printStackTrace();
System.out.println("发送消息出现异常...");
} finally {
// 7: 释放连接关闭通道
if (channel != null && channel.isOpen()) {
try {
channel.close();
} catch (Exception ex) {
ex.printStackTrace();
}
}
if (connection != null && connection.isOpen()) {
try {
connection.close();
} catch (Exception ex) {
ex.printStackTrace();
}
}
}
}
}