最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot 集成消息隊列實戰(zhàn)指南(RabbitMQ/Kafka):異步通信與解耦,落地高可靠消息傳遞

 更新時間:2026年01月20日 09:51:09   作者:小北方城市網(wǎng)  
本文介紹了SpringBoot集成RabbitMQ與Kafka的實戰(zhàn),涵蓋生產(chǎn)者、消費者、路由和可靠性保障等核心能力,幫助快速實現(xiàn)異步通信場景,解決服務(wù)耦合和流量削峰問題,感興趣的朋友跟隨小編一起看看吧

消息隊列(MQ)作為分布式系統(tǒng)的核心組件,核心價值是「異步通信、系統(tǒng)解耦、流量削峰」—— 通過消息中間件實現(xiàn)服務(wù)間的異步交互,避免服務(wù)直接調(diào)用導(dǎo)致的耦合,同時緩沖高并發(fā)流量(如秒殺、訂單峰值),保障系統(tǒng)穩(wěn)定性。主流消息隊列中,RabbitMQ 適合復(fù)雜路由、低延遲場景,Kafka 適合高吞吐、大數(shù)據(jù)場景。

本文聚焦 SpringBoot 集成 RabbitMQ 與 Kafka 的完整實戰(zhàn),嵌入可直接復(fù)用的代碼教學(xué),覆蓋生產(chǎn)者、消費者、消息路由、可靠性保障等核心能力,幫你快速落地異步通信場景,解決服務(wù)耦合、流量削峰等問題。

一、核心認(rèn)知:消息隊列的核心價值與選型

1. 核心價值

  • 系統(tǒng)解耦:服務(wù)間通過消息通信,無需感知對方存在,修改一個服務(wù)不影響其他服務(wù);
  • 異步通信:無需同步等待服務(wù)響應(yīng),發(fā)送消息后立即返回,提升接口響應(yīng)速度;
  • 流量削峰:高并發(fā)場景下,消息隊列緩沖請求,消費者按能力消費,避免下游服務(wù)被壓垮;
  • 可靠投遞:通過持久化、確認(rèn)機(jī)制,確保消息不丟失、不重復(fù)消費。

2. 選型對比(RabbitMQ vs Kafka)

特性RabbitMQKafka
吞吐量中低吞吐高吞吐(百萬級 / 秒)
延遲低延遲(毫秒級)中延遲(毫秒級)
路由能力支持復(fù)雜路由(交換機(jī))簡單路由(主題分區(qū))
可靠性強(qiáng)可靠性(確認(rèn)機(jī)制完善)可靠性可配置
適用場景訂單通知、日志告警大數(shù)據(jù)采集、秒殺削峰

二、核心實戰(zhàn)一:SpringBoot 集成 RabbitMQ(完整代碼教學(xué))

RabbitMQ 基于 AMQP 協(xié)議,核心是「交換機(jī) + 隊列 + 綁定」的路由模型,支持 Direct、Topic、Fanout 等多種交換機(jī)類型,適配復(fù)雜路由場景。

1. 環(huán)境準(zhǔn)備

(1)安裝 RabbitMQ

本地部署:Docker 命令快速啟動(推薦)

docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management
  • 5672:消息通信端口;15672:管理界面端口(訪問 http://localhost:15672,默認(rèn)賬號 guest/guest)。

(2)引入依賴(Maven)

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

(3)配置文件(application.yml)

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    virtual-host: / # 虛擬主機(jī)(默認(rèn)/)
    publisher-confirm-type: correlated # 開啟生產(chǎn)者確認(rèn)機(jī)制
    publisher-returns: true # 開啟消息回退機(jī)制
    listener:
      simple:
        acknowledge-mode: manual # 消費者手動確認(rèn)消息
        concurrency: 2 # 消費者核心線程數(shù)
        max-concurrency: 5 # 消費者最大線程數(shù)

2. 核心代碼實現(xiàn)

(1)配置類:聲明交換機(jī)、隊列、綁定關(guān)系

import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
    // 交換機(jī)名稱(Topic交換機(jī),支持模糊路由,最常用)
    public static final String TOPIC_EXCHANGE = "order_exchange";
    // 隊列名稱(訂單通知隊列)
    public static final String ORDER_QUEUE = "order_queue";
    // 路由鍵(匹配規(guī)則:order.* 匹配 order.create、order.cancel 等)
    public static final String ROUTING_KEY = "order.#";
    // 1. 聲明Topic交換機(jī)
    @Bean
    public TopicExchange topicExchange() {
        // durable=true:交換機(jī)持久化,重啟RabbitMQ不丟失
        return ExchangeBuilder.topicExchange(TOPIC_EXCHANGE).durable(true).build();
    }
    // 2. 聲明隊列
    @Bean
    public Queue orderQueue() {
        // durable=true:隊列持久化;exclusive=false:不排他;autoDelete=false:不自動刪除
        return QueueBuilder.durable(ORDER_QUEUE).build();
    }
    // 3. 綁定交換機(jī)與隊列(指定路由鍵)
    @Bean
    public Binding bindingExchangeQueue(TopicExchange topicExchange, Queue orderQueue) {
        return BindingBuilder.bind(orderQueue).to(topicExchange).with(ROUTING_KEY);
    }
}

(2)生產(chǎn)者:發(fā)送消息(含確認(rèn)機(jī)制,確保消息投遞成功)

import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.UUID;
@Component
public class OrderProducer {
    @Resource
    private RabbitTemplate rabbitTemplate;
    // 發(fā)送訂單創(chuàng)建消息
    public void sendOrderCreateMsg(Long orderId, String userId) {
        // 1. 構(gòu)建消息內(nèi)容
        String msg = String.format("用戶%s創(chuàng)建訂單:%s", userId, orderId);
        // 2. 消息ID(用于確認(rèn)機(jī)制,追蹤消息)
        CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
        // 3. 發(fā)送消息(交換機(jī)、路由鍵、消息內(nèi)容、消息ID)
        rabbitTemplate.convertAndSend(
                RabbitMQConfig.TOPIC_EXCHANGE,
                "order.create", // 具體路由鍵(匹配 order.#)
                msg,
                correlationData
        );
    }
    // 4. 生產(chǎn)者確認(rèn)回調(diào)(確認(rèn)消息是否到達(dá)交換機(jī))
    @Resource
    public void setRabbitTemplate(RabbitTemplate rabbitTemplate) {
        // 消息到達(dá)交換機(jī)回調(diào)
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (ack) {
                System.out.println("消息到達(dá)交換機(jī),消息ID:" + correlationData.getId());
            } else {
                System.out.println("消息未到達(dá)交換機(jī),原因:" + cause);
                // 消息投遞失敗,可重試或記錄日志
            }
        });
        // 消息無法路由到隊列回調(diào)(回退機(jī)制)
        rabbitTemplate.setReturnsCallback(returned -> {
            System.out.println("消息無法路由,路由鍵:" + returned.getRoutingKey() + ",原因:" + returned.getReplyText());
        });
    }
}

(3)消費者:接收消息(手動確認(rèn),確保消息消費成功)

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import java.io.IOException;
@Component
public class OrderConsumer {
    // 監(jiān)聽訂單隊列
    @RabbitListener(queues = RabbitMQConfig.ORDER_QUEUE)
    public void consumeOrderMsg(String msg, Channel channel, Message message) throws IOException {
        try {
            // 1. 處理業(yè)務(wù)邏輯(如更新訂單狀態(tài)、發(fā)送短信通知)
            System.out.println("接收訂單消息:" + msg);
            // 2. 手動確認(rèn)消息(multiple=false:只確認(rèn)當(dāng)前消息;true:確認(rèn)所有未確認(rèn)消息)
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            // 3. 消息消費失敗,拒絕消息并重回隊列(或死信隊列)
            // requeue=true:重回隊列;false:不重回隊列(需配置死信隊列處理)
            channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
            System.out.println("消息消費失敗,已重回隊列:" + msg);
        }
    }
}

3. 測試代碼(Controller 層)

import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
@RestController
public class OrderController {
    @Resource
    private OrderProducer orderProducer;
    @GetMapping("/order/create/{orderId}/{userId}")
    public String createOrder(@PathVariable Long orderId, @PathVariable String userId) {
        orderProducer.sendOrderCreateMsg(orderId, userId);
        return "訂單創(chuàng)建消息已發(fā)送";
    }
}

三、核心實戰(zhàn)二:SpringBoot 集成 Kafka(完整代碼教學(xué))

Kafka 基于發(fā)布 / 訂閱模型,核心是「主題(Topic)+ 分區(qū)(Partition)+ 消費者組(Consumer Group)」,高吞吐特性適合大數(shù)據(jù)場景。

1. 環(huán)境準(zhǔn)備

(1)安裝 Kafka

Docker 啟動單節(jié)點 Kafka(簡化版,生產(chǎn)環(huán)境需集群):

# 啟動ZooKeeper(Kafka依賴ZooKeeper管理元數(shù)據(jù))
docker run -d --name zookeeper -p 2181:2181 confluentinc/cp-zookeeper:latest
# 啟動Kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=localhost:2181 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
confluentinc/cp-kafka:latest

(2)引入依賴(Maven)

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-boot-starter-kafka</artifactId>
</dependency>

(3)配置文件(application.yml)

spring:
  kafka:
    bootstrap-servers: localhost:9092 # Kafka服務(wù)地址
    # 生產(chǎn)者配置
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      acks: 1 # 消息確認(rèn)機(jī)制(1:領(lǐng)導(dǎo)者分區(qū)確認(rèn);all:所有副本確認(rèn))
      retries: 3 # 消息發(fā)送失敗重試次數(shù)
    # 消費者配置
    consumer:
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      group-id: order-group # 消費者組ID(同一組內(nèi)消費者負(fù)載均衡消費)
      auto-offset-reset: earliest # 無偏移量時,從最早消息開始消費
      enable-auto-commit: false # 關(guān)閉自動提交偏移量,手動提交
    # 監(jiān)聽配置
    listener:
      ack-mode: manual_immediate # 手動提交偏移量

2. 核心代碼實現(xiàn)

(1)生產(chǎn)者:發(fā)送消息(異步發(fā)送,支持回調(diào))

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;
import javax.annotation.Resource;
@Component
public class KafkaOrderProducer {
    // 主題名稱(Kafka無需提前聲明,發(fā)送消息時自動創(chuàng)建)
    public static final String ORDER_TOPIC = "order_topic";
    @Resource
    private KafkaTemplate<String, String> kafkaTemplate;
    // 異步發(fā)送訂單消息
    public void sendOrderMsg(Long orderId, String userId) {
        String msg = String.format("用戶%s創(chuàng)建訂單:%s", userId, orderId);
        // 發(fā)送消息(主題、消息鍵、消息內(nèi)容)
        ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(ORDER_TOPIC, orderId.toString(), msg);
        // 回調(diào)函數(shù):處理發(fā)送結(jié)果
        future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
            @Override
            public void onSuccess(SendResult<String, String> result) {
                System.out.println("消息發(fā)送成功,分區(qū):" + result.getRecordMetadata().partition());
            }
            @Override
            public void onFailure(Throwable ex) {
                System.out.println("消息發(fā)送失敗,原因:" + ex.getMessage());
                // 失敗重試邏輯
            }
        });
    }
}

(2)消費者:接收消息(手動提交偏移量)

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
@Component
public class KafkaOrderConsumer {
    // 監(jiān)聽訂單主題,手動提交偏移量
    @KafkaListener(topics = KafkaOrderProducer.ORDER_TOPIC, groupId = "order-group")
    public void consumeOrderMsg(ConsumerRecord<String, String> record, Acknowledgment acknowledgment) {
        try {
            // 1. 獲取消息內(nèi)容
            String key = record.key();
            String msg = record.value();
            System.out.println("接收Kafka消息,訂單ID:" + key + ",內(nèi)容:" + msg);
            // 2. 處理業(yè)務(wù)邏輯
            // 3. 手動提交偏移量(確認(rèn)消息消費成功)
            acknowledgment.acknowledge();
        } catch (Exception e) {
            System.out.println("消息消費失敗:" + e.getMessage());
            // 消費失敗可記錄日志,人工介入處理
        }
    }
}

3. 測試代碼(復(fù)用上面的 OrderController,注入 KafkaOrderProducer 即可)

@Resource
private KafkaOrderProducer kafkaOrderProducer;
@GetMapping("/kafka/order/create/{orderId}/{userId}")
public String createOrderByKafka(@PathVariable Long orderId, @PathVariable String userId) {
    kafkaOrderProducer.sendOrderMsg(orderId, userId);
    return "Kafka訂單消息已發(fā)送";
}

四、消息可靠性保障(企業(yè)級實戰(zhàn)必備)

1. 消息不丟失

  • 生產(chǎn)者:開啟確認(rèn)機(jī)制(RabbitMQ 確認(rèn) + 回退,Kafka acks=all + 重試);
  • 中間件:消息持久化(RabbitMQ 交換機(jī) / 隊列持久化,Kafka 主題分區(qū)持久化);
  • 消費者:手動確認(rèn)消息,避免自動確認(rèn)導(dǎo)致消費失敗后消息丟失。

2. 消息不重復(fù)消費

  • 核心方案:基于業(yè)務(wù)唯一標(biāo)識(如訂單 ID)做冪等性處理,消費前先檢查消息是否已處理;
  • 示例代碼(消費端冪等處理):
// 基于Redis實現(xiàn)冪等(消費前檢查消息是否已處理)
public void consumeWithIdempotent(String msgId, String msg) {
    String key = "msg:processed:" + msgId;
    Boolean isProcessed = stringRedisTemplate.opsForValue().setIfAbsent(key, "1", 24, TimeUnit.HOURS);
    if (Boolean.TRUE.equals(isProcessed)) {
        // 未處理,執(zhí)行業(yè)務(wù)邏輯
        System.out.println("處理消息:" + msg);
    } else {
        // 已處理,直接跳過
        System.out.println("消息已重復(fù)消費,跳過:" + msg);
    }
}

五、避坑指南

坑點 1:RabbitMQ 消息堆積,消費者消費緩慢

表現(xiàn):隊列消息堆積過多,下游服務(wù)處理不及時;? 解決方案:增加消費者線程數(shù)(調(diào)整 concurrency/max-concurrency),拆分隊列,避免單個隊列承載過多消息。

坑點 2:Kafka 消費者組配置錯誤,導(dǎo)致重復(fù)消費

表現(xiàn):同一消費者組內(nèi)多個消費者消費同一消息;? 解決方案:確保同一主題的消費者在同一消費者組,且分區(qū)數(shù)≥消費者數(shù)(負(fù)載均衡的前提)。

坑點 3:消息確認(rèn)機(jī)制未配置,導(dǎo)致消息丟失

表現(xiàn):生產(chǎn)者發(fā)送消息后,中間件宕機(jī),消息丟失;? 解決方案:生產(chǎn)環(huán)境必須開啟生產(chǎn)者確認(rèn)機(jī)制和消息持久化,缺一不可。

六、終極總結(jié):消息隊列實戰(zhàn)的核心是「異步解耦 + 可靠傳遞」

消息隊列的本質(zhì)是「中介者」,通過它實現(xiàn)服務(wù)間的間接通信,核心價值不在于「發(fā)送消息」,而在于「安全、高效地傳遞消息」,同時解耦系統(tǒng)、緩沖流量。

核心原則總結(jié):

  1. 選型按需:復(fù)雜路由選 RabbitMQ,高吞吐大數(shù)據(jù)選 Kafka,不盲目跟風(fēng);
  2. 可靠性優(yōu)先:生產(chǎn)環(huán)境必須保障消息不丟失、不重復(fù)消費,冪等性是底線;
  3. 性能可控:合理配置消費者線程、隊列 / 分區(qū)數(shù)量,避免消息堆積或資源浪費。

記住:消息隊列不是「萬能的」,過度依賴會增加系統(tǒng)復(fù)雜度,僅在需要異步、解耦、削峰的場景使用,才能最大化其價值。

到此這篇關(guān)于SpringBoot 集成消息隊列實戰(zhàn)(RabbitMQ/Kafka):異步通信與解耦,落地高可靠消息傳遞的文章就介紹到這了,更多相關(guān)SpringBoot 集成消息隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java?JNI的高級用法示例詳解

    Java?JNI的高級用法示例詳解

    JNI是Java高級應(yīng)用中不可或缺的技術(shù)之一,它允許Java程序與本地代碼進(jìn)行交互,大大拓寬了Java應(yīng)用的范圍和性能,這篇文章主要介紹了Java?JNI高級用法的相關(guān)資料,需要的朋友可以參考下
    2025-11-11
  • Java中死鎖與活鎖的具體實現(xiàn)

    Java中死鎖與活鎖的具體實現(xiàn)

    鎖發(fā)生在不同的請求中,本文主要介紹了Java中死鎖與活鎖,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-05-05
  • Sharding-jdbc報錯:Missing the data source name:‘m0‘解決方案

    Sharding-jdbc報錯:Missing the data source 

    在使用MyBatis-plus進(jìn)行數(shù)據(jù)操作時,新增Order實體屬性后,出現(xiàn)了數(shù)據(jù)源缺失的提示錯誤,原因是因為userId屬性值使用了隨機(jī)函數(shù)生成的Long值,這與sharding-jdbc的路由規(guī)則計算不匹配,導(dǎo)致無法找到正確的數(shù)據(jù)源,通過調(diào)整userId生成邏輯
    2024-11-11
  • Nacos的單機(jī)模式啟動失敗問題及解決

    Nacos的單機(jī)模式啟動失敗問題及解決

    這篇文章主要介紹了Nacos的單機(jī)模式啟動失敗問題及解決,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-07-07
  • SpringBoot同時集成Mybatis和Mybatis-plus框架

    SpringBoot同時集成Mybatis和Mybatis-plus框架

    在實際開發(fā)中,項目里面一般都是Mybatis和Mybatis-Plus公用,但是公用有版本不兼容的問題,本文主要介紹了Spring Boot項目中同時集成Mybatis和Mybatis-plus,具有一檔的參考價值,感興趣的可以了解一下
    2024-12-12
  • maven settings.xml文件的存放及配置(包含了配置阿里云鏡像)

    maven settings.xml文件的存放及配置(包含了配置阿里云鏡像)

    本文詳細(xì)解釋了Maven中settings.xml文件的存放位置,以及用戶級別和全局級別的區(qū)別,重點介紹了localRepository、交互模式、離線模式、插件組、代理設(shè)置、服務(wù)器認(rèn)證、鏡像列表和激活profiles的使用方法,感興趣的可以了解一下
    2025-09-09
  • 基于SpringBoot使用Tika實現(xiàn)文檔解析

    基于SpringBoot使用Tika實現(xiàn)文檔解析

    Apache?Tika是開源內(nèi)容分析工具,支持多格式文本提取與元數(shù)據(jù)解析,具備語言檢測和MIME類型識別功能,適用于搜索引擎、數(shù)據(jù)分析等場景,在SpringBoot中集成需注意性能及配置問題,支持流式處理和自定義擴(kuò)展,下面介紹SpringBoot使用Tika實現(xiàn)文檔解析,感興趣的朋友一起看看吧
    2025-07-07
  • Struts2 Result 返回JSON對象詳解

    Struts2 Result 返回JSON對象詳解

    這篇文章主要講解Struts2返回JSON對象的兩種方式,講的比較詳細(xì),希望能給大家做一個參考。
    2016-06-06
  • 舉例分析Python中設(shè)計模式之外觀模式的運用

    舉例分析Python中設(shè)計模式之外觀模式的運用

    這篇文章主要介紹了Python中設(shè)計模式之外觀模式的運用,外觀模式主張以分多模塊進(jìn)行代碼管理而減少耦合,需要的朋友可以參考下
    2016-03-03
  • Spring中@PropertySource配置的用法

    Spring中@PropertySource配置的用法

    這篇文章主要介紹了Spring中@PropertySource配置的用法,@PropertySource 和 @Value
    組合使用,可以將自定義屬性文件中的屬性變量值注入到當(dāng)前類的使用@Value注解的成員變量中,需要的朋友可以參考下
    2023-11-11

最新評論

马山县| 嵊泗县| 田林县| 屏边| 花莲市| 临潭县| 陕西省| 大同县| 海阳市| 渑池县| 衡南县| 隆回县| 临澧县| 娱乐| 偏关县| 涿鹿县| 台中市| 山东省| 嘉祥县| 宣威市| 红原县| 宁津县| 内黄县| 樟树市| 牟定县| 贡觉县| 田东县| 青岛市| 虎林市| 延庆县| 花莲市| 花垣县| 威海市| 北票市| 涞水县| 无极县| 崇礼县| 乐昌市| 南开区| 胶州市| 莎车县|