带你玩转RabbitMQ的五种队列

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

                }

            }

        }

    }

}


请使用浏览器的分享功能分享到微信等