消息中间件

学习难度:❤️❤️❤️
问题:
解决

先来个简单的Java小案例,这样你就能理解了
传统:新增订单->扣减库存-》发送短信通知 (一步一步的来,要是短信发失败了,订单也不能创建了,这非常不合理!!!)
// 没有消息队列的紧耦合代码
class OrderService {
public void createOrder(Order order) {
// 1. 保存订单到数据库
orderDao.save(order);
// 2. 减库存(必须等)
inventoryService.deduct(order);
// 3. 发短信通知(必须等)
smsService.send(order.getPhone());
// 4. 发优惠券(必须等)
couponService.grant(order.getUserId());
// 所有步骤必须成功,一个失败就全部回滚!
}
}
改进加入消息队列:新建订单–》将接下来的操作交给消息中间件–》给他就立刻算成功了(不用等,发消息的执行让其他线程来异步执行)
// 使用消息队列后的代码
class OrderService {
public void createOrder(Order order) {
// 1. 保存订单(核心操作)
orderDao.save(order);
// 2. 发个消息到队列,立刻返回成功!
messageQueue.send("order.created", order);
return "下单成功!"; // 用户立刻看到
}
}
// 其他服务监听消息
class SmsService {
// 监听订单创建消息
@Listener("order.created")
public void sendSms(Order order) {
// 这里慢慢发短信,失败了重试
}
}
初学者先用Redis模拟消息队列(List数据类型)
轻量,快,简单,适合小系统,学习入门,实时通知,List (先进先出),相当于队列帮你存了一会儿,数据,让他排好队,一个一个取出来。

# 1. 饭店创建一个"接单队列"
LPUSH restaurant_orders "订单1:黄焖鸡*1"
LPUSH restaurant_orders "订单2:鱼香肉丝*2"
# 2. 查看所有待处理订单
LRANGE restaurant_orders 0 -1
# 3. 饭店从队列右边取订单处理(先进先出)
RPOP restaurant_orders
# 输出:"订单1:黄焖鸡*1"
# 4. 再看队列里还剩什么
LRANGE restaurant_orders 0 -1
# 输出:只剩订单2了
消息中间件,那么问题来了
接下我们就展开仔细探讨,按顺序我们先学Rabbit MQ因为比较全面
RabbitMQ15672
功能全面,稳定http://localhost:15672 默认的Rabbit MQ管理界面(你得先下载Docker下载镜像最快),订单,支付,通知
企业最经典的场景:创建订单通知:1.库存系统2.日志系统3.营销系统
1.打开管理界面

2.创建一个名为order_queue 的队列
3.导入客户端依赖
无论生产者/消费者都是客户端!!!!!
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.16.0</version>
</dependency>
4.生产者:发布消息
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.Channel;
public class OrderProducer {
// 队列名称
private final static String QUEUE_NAME = "order_queue";
public static void main(String[] args) throws Exception {
// 1. 创建连接工厂(相当于配置快递公司地址)
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost"); // RabbitMQ地址
factory.setUsername("guest"); // 默认用户名
factory.setPassword("guest"); // 默认密码
// 2. 建立连接(相当于打电话给快递公司)
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 3. 声明队列(如果不存在就创建)
// 参数:队列名, 是否持久化, 是否独占, 是否自动删除, 其他参数
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 4. 发送消息
String message = "订单号: 20231225001, 商品: iPhone 15, 用户: 张三";
// 参数:交换机名, 路由键, 属性, 消息体
// 这里用默认交换机,路由键=队列名
channel.basicPublish("", QUEUE_NAME, null, message.getBytes());
System.out.println(" [生产者] 已发送: " + message);
}
}
}
5.消费者:处理消息
import com.rabbitmq.client.*;
public class OrderConsumer {
private final static String QUEUE_NAME = "order_queue";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
// 声明队列(确保存在)
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
System.out.println(" [消费者] 等待订单消息…");
// 创建消费者
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [消费者] 收到新订单: " + message);
// 模拟处理订单
System.out.println(">>> 开始处理订单…");
Thread.sleep(2000); // 模拟处理时间
System.out.println(">>> 订单处理完成!");
};
// 监听队列
// 参数:队列名, 自动确认, 回调函数, 取消回调
channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> {});
}
}
6.结果
- 先运行消费者,等待接单状态(先开店后来客),在运行生产者发送消息
- 生产者—》连接RabbitMQ—-》发送消息队列
- 消费者—》连接到RabbitMQ—》监听队列—–》收到消息—-》处理
- 注意:消费者可以不在线,消息会存在消息队列中 ,等待消费者上线处理
- 精髓总结:连接工厂connectionFactory,连接connection,管道channel,管道发布/消费basicPublish/basicConsume
7.交换机
就是publish 发布消息的时候选择消息发送的–模式(比如 点对点,广播群发,)
🎯 四种交换机类型快速理解
|
Direct |
精确快递 |
"direct" |
订单路由、点对点任务 |
|
Fanout |
小区广播 |
"fanout" |
群发通知、日志收集 |
|
Topic |
智能分类 |
"topic" |
消息分类、复杂路由 |
|
Headers |
复杂筛选 |
"headers" |
特殊匹配需求(少用) |
最常用的是 Direct 和 Fanout!
经典点对点direct
public class ExchangeExample {
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection conn = factory.newConnection();
Channel channel = conn.createChannel()) {
// 1. 声明交换机(在管理网页/界面也能创建)
String EXCHANGE_NAME = "order_exchange";
String QUEUE_NAME = "order_queue";
String ROUTING_KEY = "order.created";
// 创建直连交换机
channel.exchangeDeclare(EXCHANGE_NAME, "direct", true);
// 2. 声明队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 3. 绑定队列到交换机
// 含义:把 order_queue 绑定到 order_exchange,
// 只接收 routingKey 为 "order.created" 的消息
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
// 4. 发送消息(通过交换机)
String message = "新订单消息";
channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, null, message.getBytes());
System.out.println("消息已通过交换机发送!");
}
}
}
注意:要记得声明队列 queueDeclare
8.实战订单通知系统(Fanout广播模式)
订单创建—.通知三个系统
// 生产者:发送到fanout交换机
channel.exchangeDeclare("order_fanout", "fanout", true);
channel.basicPublish("order_fanout", "", null, "新订单消息".getBytes());
// fanout交换机忽略routingKey,所以第二个参数是"",因为不是点对点
// 消费者1:库存系统
channel.queueDeclare("inventory_queue", true, false, false, null);
channel.queueBind("inventory_queue", "order_fanout", "");
// 绑定到同一个交换机
// 消费者2:日志系统
channel.queueDeclare("log_queue", true, false, false, null);
channel.queueBind("log_queue", "order_fanout", "");
// 消费者3:营销系统
channel.queueDeclare("promotion_queue", true, false, false, null);
channel.queueBind("promotion_queue", "order_fanout", "");
消息丢了怎么办?声明队列设置持久化即可。用管道channel声明队列 queueDeclare(队列名,是否持久,是否独占,是否自动删除,其他参数)
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.deliveryMode(2) // 2=持久化
.build();
channel.basicPublish("", QUEUE, props, message.getBytes());
消费者没确认消费完毕消息,?将消费者consume 自动确认关闭,改成手动确认
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.deliveryMode(2) // 2=持久化
.build();
channel.basicPublish("", QUEUE, props, message.getBytes());
记住核心思想:RabbitMQ是智能路由器,而Redis只是信箱。 理解了这个区别,你就入门了!
Kafka9092
吞吐量巨大,持久化,日志收集,大数据分析,数据的流式处理,数据被消费完成还会(按照时间几天/几月)保存,Rabbit 消息被消费就没了,消息的持久化存储,可以被多个消费者重复消费

核心概念
举例理解:topic (人,猫,鱼)
parttition(小明,小红,小兰)。(三花,梨花,豆豆),(草鱼,京腔鱼)

1.导入客户端依赖
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.4.0</version>
</dependency>
2.生产者
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class UserBehaviorProducer {
public static void main(String[] args) {
// 1. 配置生产者
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092"); // Kafka地址
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 2. 创建生产者
Producer<String, String> producer = new KafkaProducer<>(props);
// 3. 发送消息
String topic = "user_behavior";
for (int i = 1; i <= 10; i++) {
// 构建消息
String userId = "user_" + (i % 3); // 3个用户循环
String behavior = "click_product_" + i;
String message = String.format(
"{\\"user_id\\":\\"%s\\",\\"behavior\\":\\"%s\\",\\"time\\":%d}",
userId, behavior, System.currentTimeMillis()
);
// 创建ProducerRecord
ProducerRecord<String, String> record =
new ProducerRecord<>(topic, userId, message);
// 异步发送,带回调
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception e) {
if (e == null) {
System.out.printf("消息发送成功! Topic: %s, Partition: %d, Offset: %d\\n",
metadata.topic(), metadata.partition(), metadata.offset());
} else {
e.printStackTrace();
}
}
});
System.out.println("发送: " + message);
try { Thread.sleep(1000); } catch (InterruptedException e) {}
}
// 4. 关闭生产者
producer.close();
}
}
3.消费者
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class UserBehaviorConsumer {
public static void main(String[] args) {
// 1. 配置消费者
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "user-behavior-group"); // 消费者组
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("auto.offset.reset", "earliest"); // 从最早开始消费
// 2. 创建消费者
Consumer<String, String> consumer = new KafkaConsumer<>(props);
// 3. 订阅主题
consumer.subscribe(Collections.singletonList("user_behavior"));
System.out.println("开始监听用户行为…");
// 4. 持续拉取消息
try {
while (true) {
// 拉取消息,最多等待100ms
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.println("======= 收到用户行为 =======");
System.out.println("主题: " + record.topic());
System.out.println("分区: " + record.partition());
System.out.println("键: " + record.key());
System.out.println("值: " + record.value());
System.out.println("偏移量: " + record.offset());
System.out.println("时间戳: " + record.timestamp());
System.out.println("===========================\\n");
}
}
} finally {
consumer.close();
}
}
}
4.结果
先启动消费者,再启动生产者。
kafka 的优势
1. 消息持久化 – 可以"回看"
// 消费者可以从任意位置开始消费
props.put("auto.offset.reset", "earliest"); // 从最早
// 或
props.put("auto.offset.reset", "latest"); // 从最新
// 或指定offset
consumer.seek(new TopicPartition("topic", 0), 100); // 从offset=100开始
2. 消费者组 – 负载均衡
消费者组A: {Consumer1, Consumer2, Consumer3}
消费者组B: {Consumer4, Consumer5}
Topic: my-topic (3个分区)
结果:
– 组A的3个消费者每人消费1个分区
– 组B的2个消费者消费所有分区(其中一人消费2个分区)
– 两组互不影响!
3. 分区策略 – 消息路由
// 默认:轮询分区
ProducerRecord<String, String> record1 =
new ProducerRecord<>(topic, "message1"); // 轮询选择分区
// 指定分区
ProducerRecord<String, String> record2 =
new ProducerRecord<>(topic, 0, "key1", "message2"); // 发送到分区0
// 指定Key,相同Key去同一分区(保证有序!)
ProducerRecord<String, String> record3 =
new ProducerRecord<>(topic, "user123", "user action"); // user123的消息都去同一个分区
常见问题解决
问题1:消息顺序问题(生产记录发消息的时候 ,根据id 分区 partition ProducerRecord)
// 错误:不同订单的消息可能乱序
new ProducerRecord<>(topic, "order message");
// 正确:相同key的消息去同一分区,保证顺序
new ProducerRecord<>(topic, orderId, "order message");
问题2:重复消费 (消费处理完,提交的时候系统崩溃,消费前先if 判断getId 是否已经处理过)
// 原因:消费者处理完,但提交offset前崩溃
// 解决:保证幂等性
public void processMessage(Message msg) {
// 先检查是否已处理
if (isProcessed(msg.getId())) {
return; // 已处理,跳过
}
// 处理消息…
saveToDB(msg);
// 记录已处理
markAsProcessed(msg.getId());
}
问题3:消费者卡住
// 增加session超时时间
props.put("session.timeout.ms", 30000);
// 增加最大poll间隔
props.put("max.poll.interval.ms", 300000);
下篇讲这个,如果对您有帮助的话给个赞。
网址:https://blog.csdn.net/2401_85190702/article/details/157247120?spm=1001.2014.3001.5502
https://blog.csdn.net/2401_85190702/article/details/157247120?spm=1001.2014.3001.5502
RocketMQ(用的最多,国产核心学习)
分布式强(阿里巴巴),电商,金融,高并发场景





