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

一文詳解SpringBoot使用Kafka如何保證消息不丟失

 更新時間:2025年01月26日 11:31:15   作者:小信丶  
這篇文章主要為大家詳細(xì)介紹了SpringBoot使用Kafka如何保證消息不丟失的相關(guān)知識,文中的示例代碼講解詳細(xì),有需要的小伙伴可以參考一下

概述

在 Spring Boot 中使用 Kafka 時,要確保消息不丟失,主要涉及到生產(chǎn)者(Producer)、消費者(Consumer)以及 Kafka Broker 的配置和設(shè)計。

1. Spring Boot 與 Kafka 配置

Spring Boot 中使用 Kafka 時,可以通過 spring-kafka 來簡化配置和操作。以下是如何保證消息不丟

1.1 Producer 配置

Kafka 生產(chǎn)者是消息的發(fā)送方,確保消息的可靠性和不丟失需要配置

配置 Kafka 生產(chǎn)

在 application.yml 或application.properties 文件中

spring:
  kafka:
    producer:
      acks: all                  # 消息確認(rèn)策略:all表示等待所有副本確認(rèn)
      retries: 3                  # 發(fā)送失敗時的重試次數(shù)
      batch-size: 16384           # 每批發(fā)送的消息大小
      linger-ms: 1                # 消息發(fā)送的延遲時間(單位:毫秒)
      key-serializer: org.apache.kafka.common.serialization.StringSerializer   # Key序列化器
      value-serializer: org.apache.kafka.common.serialization.StringSerializer  # Value序列化器

重要配置

  • acks=all:生產(chǎn)leader 副本宕機或數(shù)據(jù)丟失的
  • retries=3:
  • linger-ms=1:生產(chǎn)者在發(fā)送消息時會有
  • batch-size=16384:設(shè)定

1.2 Kafka Producer 配置實例

通過 KafkaTemplate 發(fā)送消息時,可以通過 Java 配置來確保消息的可靠性。以下是如何

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
 
@Service
public class KafkaProducerService {
 
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
 
    private static final String TOPIC = "test_topic";
 
    public void sendMessage(String message) {
        kafkaTemplate.send(TOPIC, message);  // 發(fā)送消息
    }
}

KafkaTemplate.send() 會使用 `acks=acks=all 確保消息寫入 Kafka 集群時的可靠性

2. 消費者配置

消費者是 Kafka 的消息接收方。在消費過程中,確保消息的可靠性和不丟失需要使用合適的 acknowledgment(確認(rèn)機制)來保證消息的消費狀態(tài)。

2.1 Consumer 配置

Kafka 消費者的配置在 Spring Boot 中也可以通過 application.yml 或 application.properties 來進(jìn)行配置,確保消費者能夠在接收到消息后正確ack)消息。

spring:
  kafka:
    consumer:
      group-id: test-group        # 消費者所在的消費者組
      enable-auto-commit: false   # 自動提交消費位移,設(shè)置為false使用手動提交
      auto-offset-reset: earliest # 如果沒有偏移量(offset),從最早的位置開始消費
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer  # Key反序列化器
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer  # Value反序列化器

重要配置項解釋:

enable-auto-commit=false:禁用自動提交偏移量,手動提交偏移量可以防止消息丟失。自動提交可能會導(dǎo)

auto-offset-reset=earliest:如果消費者沒有消費過的偏earliest(最早的消息)開始消費,

group-id=test-group:消費者的組 ID,每個組

2.2 Consumer 手動提交偏移量

為了避免消息消費丟失,建議手動提交消息的偏移量。在 Spring Kafka 中,使用 Acknowledgment 來手動確

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.listener.MessageListener;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.config.ContainerProperties;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.MessageListener;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.listener.MessageListener;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
 
@Service
@EnableKafka
public class KafkaConsumerService {
 
    @KafkaListener(topics = "test_topic", groupId = "test-group")
    public void listen(ConsumerRecord<String, String> record, Acknowledgment acknowledgment) {
        // 消費消息后手動提交偏移量
        System.out.println("Received message: " + record.value());
        
        // 手動提交偏移量
        acknowledgment.acknowledge();
    }
}

acknowledgment.acknowledge():手動確認(rèn)

3. Kafka Broker 配置

Kafka Broker 配置是確保消息不丟失的關(guān)鍵部分。Kafka Broker 是 Kafka 集群的核心,它負(fù)責(zé)接收、存儲、處理和分發(fā)消息。以下是關(guān)于 Kafka Broker 配置的幾個重要方面:

3.1 副本數(shù)(Replication)與分區(qū)數(shù)(Partition)

Kafka 使用 副本(Replication) 來保證數(shù)據(jù)的高可用性和容錯性。每個主題(Topic)通常有多個 分區(qū)(Partition),每個分區(qū)都有一個 領(lǐng)導(dǎo)副本(Leader) 和多個 跟隨副本(Follower)。如果某個 Broker 或分區(qū)失敗,副本機制確保數(shù)據(jù)不會丟失,系統(tǒng)仍能正常運行。

分區(qū)數(shù):每個主題通常會有多個分區(qū),分區(qū)數(shù)的設(shè)置直接影響消息的吞吐量和并行度。更多的分區(qū)可以提高消息的處理速度,但也會增加系統(tǒng)的管理開銷。

副本數(shù):每個分區(qū)應(yīng)有多個副本,確保數(shù)據(jù)的冗余存儲。副本數(shù)設(shè)置得越高,數(shù)據(jù)丟失的可能性越小,但會占用更多的存儲資源。

關(guān)鍵配置:

# 每個主題的分區(qū)數(shù)(默認(rèn)3分區(qū))
num.partitions=3
 
# 默認(rèn)副本因子
default.replication.factor=3
 
# 最小同步副本數(shù)
min.insync.replicas=2

num.partitions:設(shè)置默認(rèn)的分區(qū)數(shù)(可以在創(chuàng)建主題時指定)。

default.replication.factor:設(shè)置每個主題的默認(rèn)副本數(shù),建議至少設(shè)置為 3。

min.insync.replicas:設(shè)置最小同步副本數(shù),確保消息只有在至少有該數(shù)量的副本成功寫入后才算成功。通常設(shè)置為 2 或更多,以保證消息的可靠性。

注意:min.insync.replicas 是為了避免某些副本沒有同步就確認(rèn)消息,這可能導(dǎo)致數(shù)據(jù)丟失。

3.2 消息持久化與日志配置

Kafka 的持久化策略和日志配置確保消息在磁盤上長期存儲。Kafka 將每個分區(qū)的消息存儲在磁盤上的日志文件中。為了避免消息丟失和提高讀寫性能,需要合理配置日志存儲相關(guān)參數(shù)。

關(guān)鍵配置:

# 消息日志的清理策略:可以選擇基于時間清理或基于大小清理
log.retention.hours=168      # 消息保留 7 天
log.retention.bytes=10737418240  # 設(shè)置日志最大存儲大小為 10 GB
 
# 消息寫入磁盤的刷新頻率
log.flush.interval.messages=10000   # 每 10000 條消息后刷新一次磁盤
log.flush.interval.ms=1000           # 每秒刷新一次磁盤
 
# 啟用壓縮
log.cleanup.policy=compact   # 啟用日志壓縮,保留最新版本的數(shù)據(jù)
 
# 消息寫入到磁盤的時延
log.segment.bytes=1073741824  # 每個日志段的大小為 1 GB

log.retention.hours:設(shè)置消息保留的時間(單位:小時),默認(rèn) 168 小時(即 7 天)。過期的消息會被刪除。

log.retention.bytes:設(shè)置日志保留的最大字節(jié)數(shù),超出大小的日志會被清理。

log.flush.interval.messages 和 log.flush.interval.ms:這些配置控制 Kafka 將消息刷新到磁盤的頻率,可以根據(jù)實際需要優(yōu)化,減少 I/O 操作。

log.cleanup.policy:設(shè)置日志清理策略??梢赃x擇 delete(基于時間或大小清理)或 compact(日志壓縮,僅保留最新版本的消息)。

log.segment.bytes:設(shè)置每個日志段的大小,當(dāng)日志達(dá)到該大小時會創(chuàng)建新的日志段。

3.3 分區(qū)副本同步與數(shù)據(jù)一致性

Kafka 的分區(qū)副本同步策略決定了數(shù)據(jù)的寫入和同步方式。為了確保消息的可靠性,需要配置副本同步機制,使數(shù)據(jù)被可靠地寫入到多個副本。

acks=all:確保生產(chǎn)者消息寫入所有副本后才確認(rèn)消息寫入成功??梢耘渲迷谏a(chǎn)者端。

min.insync.replicas:設(shè)置最小同步副本數(shù),保證數(shù)據(jù)在多個副本同步后再提交。

關(guān)鍵配置:

# 設(shè)置同步副本數(shù),確保寫入成功后至少需要此數(shù)目的副本同步
acks=all   # 等待所有副本確認(rèn)寫入
 
# 設(shè)置最小同步副本數(shù)
min.insync.replicas=2  # 確保至少有 2 個副本同步

acks=all:確保消息在所有副本成功寫入后才返回確認(rèn)。即使部分副本未同步,也不會確認(rèn)寫入。

min.insync.replicas:確保消息寫入時,至少有 min.insync.replicas 個副本同步,以避免因為單個副本失效導(dǎo)致的數(shù)據(jù)丟失。

3.4 Kafka Broker 高可用配置

在 Kafka 集群中,為了防止單點故障,Kafka 提供了高可用性配置。配置多個 Kafka Broker 和分布式部署是保證 Kafka 高可用性的基礎(chǔ)。

關(guān)鍵配置:

# Kafka Broker 啟動時的監(jiān)聽地址和端口
listeners=PLAINTEXT://localhost:9092
 
# 監(jiān)聽的內(nèi)網(wǎng)地址,通常與外網(wǎng)地址分開配置
advertised.listeners=PLAINTEXT://broker1.example.com:9092
 
# Zookeeper 配置
zookeeper.connect=localhost:2181  # 設(shè)置 Zookeeper 地址,確保 Kafka 集群管理的協(xié)調(diào)和一致性

listeners:配置 Kafka Broker 啟動時監(jiān)聽的地址和端口。

advertised.listeners:配置 Kafka Broker 向外暴露的地址和端口,消費者和生產(chǎn)者會連接到這個地址。

zookeeper.connect:設(shè)置 Zookeeper 地址,Kafka 依賴 Zookeeper 來管理集群的元數(shù)據(jù)和協(xié)調(diào)。

3.5 日志段與磁盤空間管理

Kafka 使用日志段(Log Segment)來管理消息存儲,每個分區(qū)的消息都會分布在多個日志段中,磁盤空間需要合理管理以避免磁盤溢出。

關(guān)鍵配置:

# 設(shè)置日志段的大小(單位:字節(jié))
log.segment.bytes=1073741824  # 每個日志段的大小為 1 GB
 
# 每個日志段的最大時間
log.roll.ms=604800000  # 設(shè)置為 7 天

log.segment.bytes:每個日志段的大小,達(dá)到該大小時會創(chuàng)建新的日志段。設(shè)置過大可能影響磁盤操作的效率,設(shè)置過小可能增加管理開銷。

log.roll.ms:設(shè)置日志段的滾動周期,控制日志文件分割的時間粒度。

3.6 性能優(yōu)化與磁盤 I/O 配置

Kafka 作為一個高吞吐量的分布式消息系統(tǒng),對于磁盤 I/O 和網(wǎng)絡(luò) I/O 的要求很高。合理配置 Kafka Broker 的磁盤 I/O 能提高消息的寫入和讀取速度,避免系統(tǒng)性能瓶頸。

關(guān)鍵配置:

# 設(shè)置消息的最小寫入延遲
log.min.cleanable.dirty.ratio=0.5  # 控制消息清理的閾值
 
# 設(shè)置最大允許的線程數(shù)
num.io.threads=8     # I/O 線程數(shù)量
 
# 設(shè)置 Kafka 內(nèi)存緩存大?。▎挝唬篗B)
log.buffer.size=10485760  # 默認(rèn)為 10MB

log.min.cleanable.dirty.ratio:設(shè)置日志清理的閾值,避免過多無效日志占用空間。

num.io.threads:配置 I/O 線程數(shù)。根據(jù)機器性能,適當(dāng)增加線程數(shù)以提高并發(fā)處理能力。

log.buffer.size:Kafka 寫入消息時會使用緩沖區(qū),適當(dāng)增加緩沖區(qū)大小可以提高消息的寫入性能。

3.7 Kafka 集群的監(jiān)控與故障診斷

為了確保 Kafka Broker 的高可用性和穩(wěn)定性,需要對 Kafka 集群進(jìn)行實時監(jiān)控。Kafka 提供了豐富的監(jiān)控指標(biāo),可以通過 JMX 進(jìn)行監(jiān)控,也可以結(jié)合其他工具(如 Prometheus 和 Grafana)來實現(xiàn)集群健康檢查和報警。

JMX 監(jiān)控:Kafka 提供了大量的 JMX 指標(biāo),可以用于監(jiān)控消息的吞吐量、延遲、磁盤使用情況等。

日志監(jiān)控:Kafka 的日志記錄了大量的運行時信息,定期檢查日志有助于提前發(fā)現(xiàn)潛在的問題。

4. 總結(jié)

通過合理配置和設(shè)計 Kafka 消息生產(chǎn)和消費的架構(gòu),可以有效地確保消息不丟失,并提高系統(tǒng)的可靠性和容錯性。以下是關(guān)鍵的策略總結(jié):

生產(chǎn)者配置:

  • 使用 acks=all 確保消息寫入到所有副本。
  • 啟用消息冪等性和重試機制。
  • 使用合理的 retries 和 delivery.timeout.ms 設(shè)置。

消費者配置:

  • 使用手動提交消費位移,避免消息丟失。
  • 使用去重策略和冪等性消費,確保消息不會重復(fù)消費。

Broker 配置:

  • 配置合理的副本數(shù)和分區(qū)數(shù),確保高可用性。
  • 使用 Zookeeper 確保 Kafka 集群的協(xié)調(diào)和管理。
  • 設(shè)置合理的日志清理策略,避免存儲空間耗盡。

監(jiān)控與報警:

  • 使用 JMX 監(jiān)控 Kafka 的運行狀態(tài)。
  • 結(jié)合 Spring Boot Actuator 進(jìn)行健康檢查和性能監(jiān)控。

到此這篇關(guān)于一文詳解SpringBoot使用Kafka如何保證消息不丟失的文章就介紹到這了,更多相關(guān)SpringBoot Kafka保證消息不丟失內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • mybatisplus之Wrappers.ne踩坑記錄解決

    mybatisplus之Wrappers.ne踩坑記錄解決

    這篇文章主要為大家介紹了mybatisplus之Wrappers.ne踩坑記錄解決,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-05-05
  • Maven項目部署到服務(wù)器設(shè)置訪問路徑以及配置虛擬目錄的方法

    Maven項目部署到服務(wù)器設(shè)置訪問路徑以及配置虛擬目錄的方法

    今天小編就為大家分享一篇關(guān)于Maven項目部署到服務(wù)器設(shè)置訪問路徑以及配置虛擬目錄的方法,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-02-02
  • Java面試題沖刺第十一天--集合框架篇(2)

    Java面試題沖刺第十一天--集合框架篇(2)

    這篇文章主要為大家分享了最有價值的兩道集合框架的面試題,涵蓋內(nèi)容全面,包括數(shù)據(jù)結(jié)構(gòu)和算法相關(guān)的題目、經(jīng)典面試編程題等,感興趣的小伙伴們可以參考一下
    2021-07-07
  • Java接口定義與實現(xiàn)方法分析

    Java接口定義與實現(xiàn)方法分析

    這篇文章主要介紹了Java接口定義與實現(xiàn)方法,簡單說明了接口的概念、功能,并結(jié)合實例形式分析了接口的相關(guān)定義與使用技巧,需要的朋友可以參考下
    2017-11-11
  • Java字符判斷的小例子

    Java字符判斷的小例子

    從鍵盤上輸入一個字符串,遍歷該字符串中的每個字符,若該字符為小寫字母,則輸出“此字符是小寫字母”;若為大寫字母,則輸出“此字符為大寫字母”;否則輸出“此字符不是字母”
    2013-09-09
  • SpringBoot統(tǒng)計接口調(diào)用耗時的三種方式

    SpringBoot統(tǒng)計接口調(diào)用耗時的三種方式

    在實際開發(fā)中,了解項目中接口的響應(yīng)時間是必不可少的事情,SpringBoot 項目支持監(jiān)聽接口的功能也不止一個,接下來我們分別以 AOP、ApplicationListener、Tomcat 三個方面去實現(xiàn)三種不同的監(jiān)聽接口響應(yīng)時間的操作,需要的朋友可以參考下
    2024-06-06
  • Java多線程、線程安全、線程池創(chuàng)建方式

    Java多線程、線程安全、線程池創(chuàng)建方式

    本文詳細(xì)介紹了多線程的基本概念、創(chuàng)建方式、常用方法、線程安全問題解決以及線程池的使用,通過具體案例和代碼示例,幫助讀者理解多線程在Java中的應(yīng)用和實現(xiàn),感興趣的朋友跟隨小編一起看看吧
    2026-03-03
  • springboot指定profiles啟動失敗問題及解決

    springboot指定profiles啟動失敗問題及解決

    這篇文章主要介紹了springboot指定profiles啟動失敗問題及解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-04-04
  • Java語法中Lambda表達(dá)式無法拋出異常的解決

    Java語法中Lambda表達(dá)式無法拋出異常的解決

    這篇文章主要介紹了Java語法中Lambda表達(dá)式無法拋出異常的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • SpringBoot請求參數(shù)接收控制指南分享

    SpringBoot請求參數(shù)接收控制指南分享

    這篇文章主要介紹了SpringBoot請求參數(shù)接收控制指南,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-04-04

最新評論

淮滨县| 定结县| 金秀| 凌源市| 济源市| 铜山县| 儋州市| 麦盖提县| 珠海市| 威远县| 墨脱县| 同江市| 南汇区| 靖宇县| 藁城市| 本溪| 永昌县| 鸡西市| 布尔津县| 岑溪市| 武鸣县| 蓬溪县| 临西县| 叙永县| 东乌珠穆沁旗| 九寨沟县| 扬州市| 罗平县| 钟山县| 宁明县| 巴楚县| 颍上县| 宜君县| 绥芬河市| 调兵山市| 松滋市| 德格县| 邢台市| 民县| 江永县| 晴隆县|