欢迎光临
我们一直在努力

一步一步吃透,消息中间件(Redis-List,Rabbit MQ,Rocket MQ,KafKa)最全面最详细!!!

消息中间件

学习难度:❤️❤️❤️

问题:

  • 传统每位客户下单需要,跟老板沟通,老板一个人应付不过来,消息太多。
  • 顾客先来后到,有时候搞不清楚,没有明确的秩序
  • 系统双十一订单太多!!!!!,系统容易崩溃!!
  • 解决

  • 解耦 ,消费者,不用与生产者进行沟通,只需,挑选自己喜欢的物品即可,互不影响
  • 异步处理,你下单后,可以做其他事情,异步其他线程来处理
  • 削峰填谷,让订单排队,慢慢处理(顺序)
  • 扩展,在消息队列 可以增加 消费者
  • 先来个简单的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 的队列

  • 点击Queues ->Add new queue   
  • 输入 Name:order_queue ,Durability:Durable(持久化,重启不丢消息)   
  • Add queue 
  • 3.导入客户端依赖

    无论生产者/消费者都是客户端!!!!!

    <dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.16.0</version>
    </dependency>

    4.生产者:发布消息

  • ConnectionFactory  创建连接工厂factory(消息队列的host,username,password)
  • 用工厂factory 建立新连接newConnection——-connection
  • 用connection连接  创建管道  createChannel()
  • 用管道channel声明队列   queueDeclare(队列名,是否持久,是否独占,是否自动删除,其他参数)
  • 管道基础发布  basicPublish(交换机名字,路由键=队列名字,属性null,消息体(字节))
  • 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.消费者:处理消息

  • 创建 factory  连接工厂
  • 工厂设置 Host ,,,工厂newConnetion
  • connection 连接创建管道  createChannel 
  • queueDeclare 声明队列存在
  • 创建接受回调 对象  DeliveryCallback ,打印收到订单
  • 处理订单
  • 管道基础消费 basicConsume(队列名,自动确认,deliveryCallback回调函数,取消回调)
  • 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

  • 有用管道创建  交换机并且声明类型  exchangeDeclare
  • 声明队列  queueDeclare
  • 将队列绑定到交换机   queueBind  
  • 发布消息publish
  • 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 消息被消费就没了,消息的持久化存储,可以被多个消费者重复消费 

    核心概念

  • Broker :多个broker 服务器集群(一个broker 负责一部分数据)–多个经纪人
  • Topic :消息的话题分类(user_behaviour,order_events)—-话题分类
  • Partition:核心 分区 在 用户行为日志Topic 中 根据 用户ID 分多个 ——区域 
  • Producer:生产者,消息发到 Topic 话题
  • Consumer:消费者,从topic 拉取消息
  • 举例理解:topic  (人,猫,鱼)

                     parttition(小明,小红,小兰)。(三花,梨花,豆豆),(草鱼,京腔鱼)

    1.导入客户端依赖

    <dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.4.0</version>
    </dependency>

    2.生产者

  • 配置Properties属性   put 设置 kafka 网址,key ,value 
  • KafkaProducer 新建生产者带上 配置属性
  • 创建生产记录  ProducerRecord ,
  • 发送消息 producer.send 带 Callback 回调   metadata.topic(), metadata.partition(), metadata.offset());
  • 关闭生产者 producer.close();
  • 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.消费者

  • properties 配置消费者
  • 创建消费者  KafKaConsumer
  • consumer.subscribe  订阅主题 topic
  • ConsumerRecords     consumer.poll   持续拉取消息
  • consumer.close()  关闭消费者
  • 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.5502https://blog.csdn.net/2401_85190702/article/details/157247120?spm=1001.2014.3001.5502

    RocketMQ(用的最多,国产核心学习)

    分布式强(阿里巴巴),电商,金融,高并发场景

    赞(0)
    未经允许不得转载:171主机测评 » 一步一步吃透,消息中间件(Redis-List,Rabbit MQ,Rocket MQ,KafKa)最全面最详细!!!
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址