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

關(guān)于Kafka消費者訂閱方式

 更新時間:2022年05月05日 11:16:30   作者:芒果無憂  
這篇文章主要介紹了關(guān)于Kafka消費者訂閱方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教

Kafka消費者訂閱方式

Kafka為消費者提供了三種類型的訂閱消費方式:訂閱主題集合、正則表達式訂閱主題、訂閱指定主題的分區(qū)集合。三種方式只能使用其中一種。

1.指定主題消費

一個消費者可以使用KafkaConsumer提供的subscribe()方法訂閱一個或多個主題,訂閱主題集合和正則表達式訂閱主題都使用此方法實現(xiàn)的。下面兩種方式都可以訂閱topic_1120主題。

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
// 訂閱主題
consumer.subscribe(Collections.singletonList("topic_1120"));
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
//正則表達式.*代表后續(xù)0個或者多個任意字符。
consumer.subscribe(Pattern.compile("topic.*"));

訂閱主題在源碼中由4個方法重載實現(xiàn),其中兩個帶listener的方法是可以自定義Rebalance重平衡的監(jiān)聽類。

@Override
public void subscribe(Collection<String> topics) {
? ? subscribe(topics, new NoOpConsumerRebalanceListener());
}
@Override
public void subscribe(Collection<String> topics, ConsumerRebalanceListener listener) {
? ? //省略源碼
}
@Override
public void subscribe(Pattern pattern) {
? ? subscribe(pattern, new NoOpConsumerRebalanceListener());
}
@Override
public void subscribe(Pattern pattern, ConsumerRebalanceListener listener) {
? ? //省略源碼
}

2.指定分區(qū)消費

消費者指定分區(qū)消費是通過KafkaConsumer提供的assign()方法實現(xiàn)的,assign()方法入?yún)镃ollection, 其中TopicPartition有2個屬性, topic和partition, 分區(qū)從0開始編號。使用assign()方法訂閱指定主題test_1120分區(qū)0的消息。

/訂閱指定分區(qū)
consumer.assign(Collections.singleton(new TopicPartition("topic_1120", 0)));

3.取消訂閱

取消訂閱調(diào)用unsubscribe()方法。

consumer.unsubscribe();

小結(jié):subscribe()具有自動重平衡的功能,來實現(xiàn)消費負載均衡和故障自動轉(zhuǎn)移,而assign()不具備這種功能。

Kafka概述

定義

Kafka是一個分布式的基于發(fā)布/訂閱模式的消息隊列,主要應(yīng)用于大數(shù)據(jù)實時處理領(lǐng)域。

消息隊列

1.傳統(tǒng)消息隊列的應(yīng)用場景

使用消息隊列的好處

1)解耦

允許你獨立的擴展或修改兩邊的處理過程,只要確保它們遵守同樣的接口約束。

2)可恢復(fù)性

系統(tǒng)的一部分組件失效時,不會影響到整個系統(tǒng)。消息隊列降低了進程間的耦合度,所 以即使一個處理消息的進程掛掉,加入隊列中的消息仍然可以在系統(tǒng)恢復(fù)后被處理。

3)緩沖

有助于控制和優(yōu)化數(shù)據(jù)流經(jīng)過系統(tǒng)的速度,解決生產(chǎn)消息和消費消息的處理速度不一致的情況。

4)靈活性 & 峰值處理能力

在訪問量劇增的情況下,應(yīng)用仍然需要繼續(xù)發(fā)揮作用,但是這樣的突發(fā)流量并不常見。 如果為以能處理這類峰值訪問為標準來投入資源隨時待命無疑是巨大的浪費。使用消息隊列 能夠使關(guān)鍵組件頂住突發(fā)的訪問壓力,而不會因為突發(fā)的超負荷的請求而完全崩潰。

5)異步通信

很多時候,用戶不想也不需要立即處理消息。消息隊列提供了異步處理機制,允許用戶 把一個消息放入隊列,但并不立即處理它。想向隊列中放入多少消息就放多少,然后在需要 的時候再去處理它們。

2.消息隊列的兩種模式

(1)點對點模式(一對一,消費者主動拉取數(shù)據(jù),消息收到后消息清除)

消息生產(chǎn)者生產(chǎn)消息發(fā)送到Queue中,然后消息消費者從Queue中取出并且消費消息。 消息被消費以后,queue 中不再有存儲,所以消息消費者不可能消費到已經(jīng)被消費的消息。 Queue 支持存在多個消費者,但是對一個消息而言,只會有一個消費者可以消費。

(2)發(fā)布/訂閱模式(一對多,消費者消費數(shù)據(jù)之后不會清除消息)

消息生產(chǎn)者(發(fā)布)將消息發(fā)布到 topic 中,同時有多個消息消費者(訂閱)消費該消 息。和點對點方式不同,發(fā)布到 topic 的消息會被所有訂閱者消費。

Kafka 基礎(chǔ)架構(gòu)

  • Producer :消息生產(chǎn)者,就是向 kafka broker 發(fā)消息的客戶端;
  • Consumer :消息消費者,向 kafka broker 取消息的客戶端;
  • Consumer Group (CG):消費者組,由多個 consumer 組成。消費者組內(nèi)每個消費者負 責消費不同分區(qū)的數(shù)據(jù),一個分區(qū)只能由一個組內(nèi)消費者消費;消費者組之間互不影響。所 有的消費者都屬于某個消費者組,即消費者組是邏輯上的一個訂閱者。
  • Broker :一臺 kafka 服務(wù)器就是一個 broker。一個集群由多個 broker 組成。一個 broker 可以容納多個 topic。
  • Topic :可以理解為一個隊列,生產(chǎn)者和消費者面向的都是一個 topic;
  • Partition:為了實現(xiàn)擴展性,一個非常大的 topic 可以分布到多個 broker(即服務(wù)器)上, 一個 topic 可以分為多個 partition,每個 partition 是一個有序的隊列;
  • Replica:副本,為保證集群中的某個節(jié)點發(fā)生故障時,該節(jié)點上的 partition 數(shù)據(jù)不丟失, 且 kafka 仍然能夠繼續(xù)工作,kafka 提供了副本機制,一個 topic 的每個分區(qū)都有若干個副本, 一個 leader 和若干個 follower。
  • leader:每個分區(qū)多個副本的“主”,生產(chǎn)者發(fā)送數(shù)據(jù)的對象,以及消費者消費數(shù)據(jù)的對 象都是 leader。
  • follower:每個分區(qū)多個副本中的“從”,實時從 leader 中同步數(shù)據(jù),保持和 leader 數(shù)據(jù) 的同步。leader 發(fā)生故障時,某個 follower 會成為新的 follower。

以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • 解決Idea的選擇文件后定位瞄準器"Select Opened File"的功能不見了

    解決Idea的選擇文件后定位瞄準器"Select Opened File"的功能

    使用IntelliJ IDEA時,可能會發(fā)現(xiàn)"SelectOpenedFile"功能不見了,這個功能允許用戶快速定位到當前打開文件的位置,若要找回此功能,只需在IDEA的標題欄上右鍵,然后選擇"Always Select Opened File",這樣就可以重新啟用這個便捷的功能
    2024-11-11
  • 通過JWT來解決登錄認證問題的方案

    通過JWT來解決登錄認證問題的方案

    Json web token (JWT),是為了在網(wǎng)絡(luò)應(yīng)用環(huán)境間傳遞聲明而執(zhí)行的一種基于JSON的開放標準((RFC7519),該token被設(shè)計為緊湊且安全的,特別適用于分布式站點的單點登錄(SSO)場景,本文給大家介紹了如何通過 JWT 來解決登錄認證問題,需要的朋友可以參考下
    2024-12-12
  • 詳解Java8新特性之interface中的static方法和default方法

    詳解Java8新特性之interface中的static方法和default方法

    這篇文章主要介紹了Java8新特性之interface中的static方法和default方法,非常不錯,具有一定的參考借鑒價值,需要的朋友可以參考下
    2018-08-08
  • 關(guān)于jvm的垃圾回收器以及觸發(fā)full gc的場景

    關(guān)于jvm的垃圾回收器以及觸發(fā)full gc的場景

    這篇文章主要介紹了關(guān)于jvm的垃圾回收器以及觸發(fā)full gc的場景,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-04-04
  • MyBatis 延遲加載、一級緩存、二級緩存(詳解)

    MyBatis 延遲加載、一級緩存、二級緩存(詳解)

    下面小編就為大家?guī)硪黄狹yBatis 延遲加載、一級緩存、二級緩存(詳解)。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-08-08
  • Spring實現(xiàn)定時任務(wù)的幾種方式總結(jié)

    Spring實現(xiàn)定時任務(wù)的幾種方式總結(jié)

    Spring Task 是 Spring 框架提供的一種任務(wù)調(diào)度和異步處理的解決方案,可以按照約定的時間自動執(zhí)行某個代碼邏輯它可以幫助開發(fā)者在 Spring 應(yīng)用中輕松地實現(xiàn)定時任務(wù)、異步任務(wù)等功能,提高應(yīng)用的效率和可維護性,需要的朋友可以參考下本文
    2024-07-07
  • Java實現(xiàn)飛機大戰(zhàn)-II游戲詳解

    Java實現(xiàn)飛機大戰(zhàn)-II游戲詳解

    《飛機大戰(zhàn)-II》是一款融合了街機、競技等多種元素的經(jīng)典射擊手游。游戲是用java語言實現(xiàn),采用了swing技術(shù)進行了界面化處理,感興趣的可以了解一下
    2022-02-02
  • SpringBoot中注冊Bean的方式總結(jié)

    SpringBoot中注冊Bean的方式總結(jié)

    這篇文章主要介紹了SpringBoot中注冊Bean的方式總結(jié),@ComponentScan + @Componet相關(guān)注解,@Bean,@Import和spring.factories這四種方式,文中代碼示例給大家介紹的非常詳細,需要的朋友可以參考下
    2024-04-04
  • 使用MyBatis-Generator如何自動生成映射文件

    使用MyBatis-Generator如何自動生成映射文件

    這篇文章主要介紹了使用MyBatis-Generator如何自動生成映射文件,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • Java設(shè)計模式之適配器模式簡介

    Java設(shè)計模式之適配器模式簡介

    這篇文章主要介紹了Java設(shè)計模式之適配器模式,需要的朋友可以參考下
    2014-07-07

最新評論

德州市| 大化| 左权县| 郁南县| 洛阳市| 芜湖县| 聂拉木县| 青岛市| 观塘区| 崇礼县| 河南省| 林芝县| 凤山市| 旅游| 新田县| 民丰县| 清徐县| 西青区| 方山县| 固安县| 哈巴河县| 凌海市| 栾城县| 色达县| 商水县| 克东县| 琼中| 华宁县| 尉氏县| 泸西县| 无棣县| 临高县| 崇仁县| 丰宁| 临江市| 泾川县| 嘉义县| 孟州市| 西安市| 东丰县| 淳安县|