RabbitMQ簡單Java示例——生產者和消費者


添加Maven依賴:

使用rabbitmq-client的最新Maven坐標:

<!-- https://mvnrepository.com/artifact/com.rabbitmq/amqp-client -->
<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.3.0</version>
</dependency>

添加賬戶

默認情況下,訪問RabbitMQ服務的用戶名和密碼都是“guest”,這個賬號有限制,默認只能通過本地網絡(如localhost)訪問,遠程網絡訪問受限,所以在實現生產和消費消息之前,需要另外添加一個用戶,並設置相應的訪問權限。

添加新用戶,用戶名為“zifeiy”,密碼為“passwd”:

C:\Users\zifeiy>rabbitmqctl add_user zifeiy passwd
Adding user "zifeiy" ...

為zifeiy用戶設置所有權限:

C:\Users\zifeiy>rabbitmqctl set_permissions -p / zifeiy ".*" ".*" ".*"
Setting permissions for user "zifeiy" in vhost "/" ...

設置用戶zifeiy為管理員角色:

C:\Users\zifeiy>rabbitmqctl set_user_tags zifeiy administrator
Setting tags for user "zifeiy" to [administrator] ...

計算機的世界是從“Hello World!”開始的,這里我們也沿用慣例,首先生產者發送一條消息”Hello World!“至RabbitMQ中,之后由消費者消費。
下面先演示生產者客戶端的代碼,然后再演示消費者客戶端的代碼。

生產者客戶端代碼

package com.zifeiy.springtest.rabbitmq;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;

public class RabbitProducer {
	private static final String EXCHANGE_NAME = "exchange_demo";
	private static final String ROUTING_KEY = "routingkey_demo";
	private static final String QUEUE_NAME = "queue_demo";
	private static final String IP_ADDRESS = "127.0.0.1";
	private static final int PORT = 5672;	// RabbitMQ服務端默認端口號為5672
	
	public static void main(String[] args) throws IOException, TimeoutException {
		ConnectionFactory factory = new ConnectionFactory();
		factory.setHost(IP_ADDRESS);
		factory.setPort(PORT);
		factory.setUsername("zifeiy");
		factory.setPassword("passwd");
		Connection connection = factory.newConnection();	// 建立連接
		Channel channel = connection.createChannel();		// 創建信道
		// 創建一個type="direct"、持久化的、非自動刪除的交換器
		channel.exchangeDeclare(EXCHANGE_NAME, "direct", true, false, null);
		// 創建一個持久化、非排他的、非自動刪除的隊列
		channel.queueDeclare(QUEUE_NAME, true, false, false, null);
		// 將交換器和隊列通過路由綁定
		channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
		// 發送一條持久化的消息:hello world!
		String message = "hello,world!";
		channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, 
					MessageProperties.PERSISTENT_TEXT_PLAIN, 
					message.getBytes());
		// 關閉資源
		channel.close();
		connection.close();
	}
}

運行。

消費者客戶端代碼

package com.zifeiy.springtest.rabbitmq;

import java.io.IOException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Address;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
import com.rabbitmq.client.Connection;

public class RabbitConsumer {
	private static final String QUEUE_NAME = "queue_demo";
	private static final String IP_ADDRESS = "127.0.0.1";
	private static final int PORT = 5672;
	
	public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
		Address[] addresses = new Address[] {
				new Address(IP_ADDRESS, PORT)
		};
		ConnectionFactory factory = new ConnectionFactory();
		factory.setUsername("zifeiy");
		factory.setPassword("passwd");
		// 這里的連接方式與生產者的demo略有不同,注意區分
		Connection connection = factory.newConnection(addresses);	// 創建連接
		final Channel channel = connection.createChannel();	// 創建信道
		channel.basicQos(64); 	// 設置客戶端最多接受未被ack的消息的個數
		Consumer consumer = new DefaultConsumer(channel) {
			@Override
			public void handleDelivery(String consumerTag, Envelope envelope, 
							AMQP.BasicProperties properties, byte[] body) throws IOException {
				System.out.println("recv message: " + new String(body));
				try {
					TimeUnit.SECONDS.sleep(1);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				channel.basicAck(envelope.getDeliveryTag(), false);
			}
		};
		channel.basicConsume(QUEUE_NAME, consumer);
		// 等待回調函數執行完畢后,關閉資源
		TimeUnit.SECONDS.sleep(5);
		channel.close();
		connection.close();
	}
}

運行,命令行輸出如下:

recv message: hello,world!


免責聲明!

本站轉載的文章為個人學習借鑒使用,本站對版權不負任何法律責任。如果侵犯了您的隱私權益,請聯系本站郵箱yoyou2525@163.com刪除。



 
粵ICP備18138465號   © 2018-2025 CODEPRJ.COM