rabbitmq發延時訊息以及通過一個exchange發到不同的queue
阿新 • • 發佈:2019-01-01
public static void main(String[] args) throws Exception { producer(1); producer(2); producer(3); } private static void producer(int i) throws Exception{ ConnectionFactory factory = new ConnectionFactory(); factory.setHost("127.0.0.1"); factory.setUsername("admin"); factory.setPassword("123456"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); HashMap<String, Object> arguments = new HashMap<String, Object>(); arguments.put("x-dead-letter-exchange", "amq.direct"); arguments.put("x-dead-letter-routing-key", "message_ttl_routingKey" +i); channel.queueDeclare("delay_queue"+i, true, false, false, arguments); channel.queueDeclare(queue_name+i, true, false, false, null); channel.queueBind(queue_name+i, "amq.direct", "message_ttl_routingKey"+i); String message = i+":hello world!" + System.currentTimeMillis(); AMQP.BasicProperties.Builder builder = new AMQP.BasicProperties.Builder(); // 永續性 non-persistent (1) or persistent (2)
//expiration設定延時時間
AMQP.BasicProperties properties = builder.expiration("0").deliveryMode(2).build(); // routingKey =delay_queue 進行轉發 // RabbitMQ預設有一個exchange,叫default exchange,它用一個空字串表示,它是direct exchange型別// ,任何發往這個exchange的訊息都會被路由到routing key的名字對應的佇列上,如果沒有對應的佇列,則訊息會被丟棄 channel.basicPublish("", "delay_queue"+i/*發到 delay_queue */, properties, message.getBytes()); System.out.println("sent message: " + message + ",date:" + System.currentTimeMillis()); // 關閉頻道和連線 channel.close(); connection.close(); }private static String queue_name = "message_ttl_queue1"; public static void main(String[] args) throws Exception { // producer(); consumer(); } private static void consumer() throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("127.0.0.1"); factory.setUsername("admin"); factory.setPassword("123456"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); HashMap<String, Object> arguments = new HashMap<String, Object>(); arguments.put("x-dead-letter-exchange", "amq.direct"); arguments.put("x-dead-letter-routing-key", "message_ttl_routingKey1"); channel.queueDeclare("delay_queue1", true, false, false, arguments); // 宣告佇列 channel.queueDeclare(queue_name, true, false, false, null); // 繫結路由 channel.queueBind(queue_name, "amq.direct", "message_ttl_routingKey1"); QueueingConsumer consumer = new QueueingConsumer(channel); // 指定消費佇列 channel.basicConsume(queue_name, true, consumer); while (true) { // nextDelivery是一個阻塞方法(內部實現其實是阻塞佇列的take方法) QueueingConsumer.Delivery delivery = consumer.nextDelivery(); String message = new String(delivery.getBody()); System.out.println("received message:" + message + ",date:" + System.currentTimeMillis()); } }
channel.basicPublish("", "delay_queue"+i, properties, message.getBytes());
先將訊息發到 delay_queue ,在properties設定的expiration超時之後,message_ttl_routingKey將訊息路由到繫結的queue