RabbitMQ——helloworld
添加maven依赖
<dependency> <groupId>com.rabbitmqgroupId> <artifactId>amqp-clientartifactId> <version>5.8.0version> dependency>
通过下列函数,获取消息队列的连接
public class ConnectionUtil { private static Logger logger = LoggerFactory.getLogger(ConnectionUtil.class); public static Connection getConnection() { try { Connection connection = null; ConnectionFactory factory = new ConnectionFactory(); //设置服务端地址(域名地址/ip) factory.setHost("127.0.0.1"); //设置服务器端口号,最易犯的错误就是填写为15672(网页端口),正确为5672 factory.setPort(5672); //设置虚拟主机(相当于数据库中的库) factory.setVirtualHost("/"); //设置用户名 factory.setUsername("root"); //设置密码 factory.setPassword("root"); connection = factory.newConnection(); return connection; } catch (Exception e) { e.printStackTrace(); return null; } } public static void main(String[] args) { getConnection(); } }
消息发送方
package cn.seaboot.admin.rabbit.test; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import java.io.IOException; import java.nio.charset.Charset; import java.util.concurrent.TimeoutException; /** * @author Mr.css * @date 2020-11-12 19:30 */ public class Send { private static final String QUEUE_NAME = "queue_name"; public static void main(String[] args) { try { Connection connection = ConnectionUtil.getConnection(); Channel channel = connection.createChannel(); channel.queueDeclare(QUEUE_NAME, false, false, false, null); int cnt = 10; while(cnt --> 0){ String message = "This is simple queue:" + cnt; //发送消息 channel.basicPublish("", QUEUE_NAME, null, message.getBytes(Charset.defaultCharset())); System.out.println("[send]:" + message); } channel.close(); connection.close(); } catch (IOException | TimeoutException e) { e.printStackTrace(); } } }
消息消费者
/** * @author Mr.css * @date 2020-11-12 19:31 */ public class Receive { private static final String QUEUE_NAME = "queue_name"; public static void main(String[] args) { try { Connection connection = ConnectionUtil.getConnection(); Channel channel = connection.createChannel(); channel.queueDeclare(QUEUE_NAME, false, false, false, null); DefaultConsumer consumer = new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) { String message = new String(body, StandardCharsets.UTF_8); System.out.println("[Receive]:" + message); try { Thread.sleep(500); } catch (InterruptedException e) { e.printStackTrace(); } } }; channel.basicConsume(QUEUE_NAME, true, consumer); } catch (IOException | ShutdownSignalException | ConsumerCancelledException e) { e.printStackTrace(); } } }