import com.rabbitmq.client.*;
public class PersistentProducer {
private final static String QUEUE_NAME = "persistent_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(QUEUE_NAME, true, false, false, null);
// 启用生产者确认
channel.confirmSelect();
String message = "Persistent message with producer confirm!";
channel.basicPublish("", QUEUE_NAME,
MessageProperties.PERSISTENT_TEXT_PLAIN,
message.getBytes());
// 检查消息是否成功发送
if (channel.waitForConfirms()) {
System.out.println("Message sent successfully!");
} else {
System.out.println("Message failed to send!");
}
}
}
}
import com.rabbitmq.client.*;
public class DirectExchangeProducer {
private final static String EXCHANGE_NAME = "direct_logs";
private final static String QUEUE_NAME = "persistent_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(EXCHANGE_NAME, "direct");
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, "info");
String message = "Hello, Direct Exchange!";
channel.basicPublish(EXCHANGE_NAME, "info",
MessageProperties.PERSISTENT_TEXT_PLAIN,
message.getBytes());
System.out.println("Sent: " + message);
}
}
}
public class DeadLetterConsumer {
private final static String DLX_QUEUE = "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.queueDeclare(DLX_QUEUE, true, false, false, null);
// 创建消费者回调
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println("Dead Letter Queue Received: " + message);
};
// 设置死信队列消费者
channel.basicConsume(DLX_QUEUE, true, deliverCallback, consumerTag -> {});
}
}
}
import com.rabbitmq.client.*;
public class AckConsumer {
private final static String QUEUE_NAME = "persistent_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(QUEUE_NAME, true, false, false, null);
// 创建一个消费者回调
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println("Received: " + message);
try {
// 模拟消息处理
if (message.contains("error")) {
throw new Exception("Error while processing message");
}
// 手动确认消息
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
System.out.println("Message processed and acknowledged");
} catch (Exception e) {
// 如果消息处理失败,可以将消息重新放回队列
System.out.println("Error processing message, requeueing: " + e.getMessage());
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
}
};
// 设置手动确认
channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {});
}
}
}
参与评论
手机查看
返回顶部