RabbitMQ 学习示例(Java)
学习 BI 项目的时候感觉鱼哥写的例子有点乱,自己写了一些
Hello World
/** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:20:00 */ public class SingleProducer { private final static String QUEUE_NAME = "hello"; public static void main(String[] argv) throws Exception { // 创建连接 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); // factory.setUsername("guest"); // factory.setPassword("guest"); // 创建连接 try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // 创建队列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); String message = "Hello World!"; channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } /** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:28:00 */ public class SingleConsumer { private final static String QUEUE_NAME = "hello"; public static void main(String[] argv) throws Exception { // 创建连接 // 参数必须与创建时的相同 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 创建队列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); // 定义消息处理方法 DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x] Received '" + message + "'"); }; // 开启消费监听 channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> { }); } }
Work Queues
/** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:20:00 */ public class MultiProducer { private static final String TASK_QUEUE_NAME = "task_queue"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null); Scanner scanner = new Scanner(System.in); while (scanner.hasNext()) { String message = scanner.nextLine(); // 消息持久化 // MessageProperties.PERSISTENT_TEXT_PLAIN channel.basicPublish("", TASK_QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } } /** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:28:00 */ public class MultiConsumer { private static final String TASK_QUEUE_NAME = "task_queue"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); final Connection connection = factory.newConnection(); // 创建多个消费者 for (int i = 0; i < 3; i++) { final Channel channel = connection.createChannel(); // 队列持久化 // durable: true channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); channel.basicQos(1); int finalI = i; DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [" + finalI + "] Received '" + message + "'"); try { // doWork(message); channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } catch (Exception e) { channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } finally { System.out.println(" [" + finalI + "] Done"); channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } }; // 开启消费监听 // autoAck: false 不自动确认 channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> { }); } } }
交换机
fanout
public class FanoutProducer { private static final String EXCHANGE_NAME = "fanout-exchange"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); Scanner scanner = new Scanner(System.in); while (scanner.hasNext()) { String message = scanner.nextLine(); channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } } public class FanoutConsumer { private static final String EXCHANGE_NAME = "fanout-exchange"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); for (int i = 0; i < 3; i++) { Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); String queueName = "queue-fanout-exchange:" + i; channel.queueDeclare(queueName, true, false, false, null); channel.queueBind(queueName, EXCHANGE_NAME, ""); System.out.println(" [" + i + "] Waiting for messages. To exit press CTRL+C"); int finalI = i; DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [" + finalI + "] Received '" + message + "'"); }; channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { }); } } }
direct
public class DirectProducer { private static final String EXCHANGE_NAME = "direct_exchange"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); Scanner scanner = new Scanner(System.in); while (scanner.hasNext()) { String inputs = scanner.nextLine(); String[] splits = inputs.split(" "); if (splits.length < 1) { continue; } String message = splits[0]; String routingKey = splits[1]; channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + routingKey + "':'" + message + "'"); } } } } public class DirectConsumer { private static final String EXCHANGE_NAME = "direct_exchange"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); for (int i = 0; i < 3; i++) { Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); String queueName = "queue-direct-exchange:" + i; channel.queueDeclare(queueName, true, false, false, null); channel.queueBind(queueName, EXCHANGE_NAME, "route-" + i); System.out.println(" [" + i + "] Waiting for messages. To exit press CTRL+C"); int finalI = i; DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [" + finalI + "] Received '" + delivery.getEnvelope().getRoutingKey() + "':'" + message + "'"); }; channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { }); } } }
topic
public class TopicProducer { private static final String EXCHANGE_NAME = "topic_exchange"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC); Scanner scanner = new Scanner(System.in); while (scanner.hasNext()) { String inputs = scanner.nextLine(); String[] splits = inputs.split(" "); if (splits.length < 1) { continue; } String message = splits[0]; String routingKey = splits[1]; channel.basicPublish(EXCHANGE_NAME, routingKey, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + routingKey + "':'" + message + "'"); } } } } public class TopicConsumer { private static final String EXCHANGE_NAME = "topic_exchange"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); for (int i = 0; i < 3; i++) { for (int j = 0; j < 3; j++) { Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC); String queueName = "queue-topic-exchange:" + i + ":" + j; channel.queueDeclare(queueName, true, false, false, null); // channel.queueBind(queueName, EXCHANGE_NAME, "route." + i + ".*"); channel.queueBind(queueName, EXCHANGE_NAME, "route." + i + ".#"); System.out.println(" [" + i + ":" + j + "] Waiting for messages. To exit press CTRL+C"); int finalI = i; int finalJ = j; DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [" + finalI + ":" + finalJ + "] Received '" + delivery.getEnvelope().getRoutingKey() + "':'" + message + "'"); }; channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { }); } } } }
消息过期机制
/** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:20:00 */ public class TtlProducer { private final static String QUEUE_NAME = "ttl_queue"; public static void main(String[] argv) throws Exception { // 创建连接 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); // factory.setUsername("guest"); // factory.setPassword("guest"); // 创建连接 try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // 创建队列 Map<String, Object> args = new HashMap<>(); // 队列过期时间设置 args.put("x-message-ttl", 60000); channel.queueDeclare(QUEUE_NAME, false, false, false, args); String message = "Hello World!"; // 消息过期时间配置 AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder() .expiration("10000") .build(); channel.basicPublish("", QUEUE_NAME, properties, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } /** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:28:00 */ public class TtlConsumer { private final static String QUEUE_NAME = "ttl_queue"; public static void main(String[] argv) throws Exception { // 创建连接 // 参数必须与创建时的相同 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 创建队列 Map<String, Object> args = new HashMap<>(); // 队列过期时间设置 args.put("x-message-ttl", 60000); channel.queueDeclare(QUEUE_NAME, false, false, false, args); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); // 定义消息处理方法 DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x] Received '" + message + "'"); }; // 开启消费监听 // 不自动确认 channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> { }); } }
死信队列/交换机
/** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:20:00 */ public class DeadLetterProducer { private final static String EXCHANGE_NAME = "exchange_dlx_demo"; private final static String ROUTING_KEY = "routing_key_dlx"; private final static String QUEUE_NAME = "queue_dlx_demo"; // 私信队列 private final static String DLX_EXCHANGE_NAME = "dlx_exchange"; // 死信交换机 private final static String DLX_QUEUE_NAME = "dlx_queue"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // 声明死信交换机 channel.exchangeDeclare(DLX_EXCHANGE_NAME, "direct"); // 声明死信队列 channel.queueDeclare(DLX_QUEUE_NAME, true, false, false, null); channel.queueBind(DLX_QUEUE_NAME, DLX_EXCHANGE_NAME, ROUTING_KEY); // 声明主交换机 channel.exchangeDeclare(EXCHANGE_NAME, "direct"); Map<String, Object> args = new HashMap<>(); // 绑定死信交换机及死信路由 args.put("x-dead-letter-exchange", DLX_EXCHANGE_NAME); args.put("x-dead-letter-routing-key", ROUTING_KEY); // 声明主交换机 channel.queueDeclare(QUEUE_NAME, true, false, false, args); channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY); Scanner scanner = new Scanner(System.in); while (scanner.hasNext()) { String message = scanner.nextLine(); channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } } /** * @author <a href="https://github.com/cutepikachu-cn">笨蛋皮卡丘</a> * @version 1.0 * @since 2024-05-22 15:28:00 */ public class DeadLetterConsumer { private final static String QUEUE_NAME = "queue_dlx_demo"; // 死信队列 private final static String DLX_QUEUE_NAME = "dlx_queue"; public static void main(String[] argv) throws IOException, TimeoutException { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 为主队列创建消费者 DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x] Received '" + message + "'"); if (RandomUtil.randomBoolean()) { // 确认 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } else { // 拒绝消息并将其重新排队到 DLX channel.basicReject(delivery.getEnvelope().getDeliveryTag(), false); } }; channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> { }); // 为死信队列创建消费者 Channel dlxChannel = connection.createChannel(); DeliverCallback dlxDeliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [DLX] Received '" + message + "'"); }; dlxChannel.basicConsume(DLX_QUEUE_NAME, true, dlxDeliverCallback, consumerTag -> { }); } }
