主頁 >  其他 > 聊聊 Kafka:Producer 原始碼決議

聊聊 Kafka:Producer 原始碼決議

2021-09-10 07:58:51 其他

一、前言

前面幾篇我們講了關于 Kafka 的基礎架構以及搭建,從這篇開始我們就來原始碼分析一波,我們這用的 Kafka 版本是 2.7.0,其 Client 端是由 Java 實作,Server 端是由 Scala 來實作的,在使用 Kafka 時,Client 是用戶最先接觸到的部分,因此,我們從 Client 端開始,會先從 Producer 端開始,今天我們就來對 Producer 原始碼決議一番,

二、Producer 使用

首先我們先通過一段代碼來展示 KafkaProducer 的使用方法,在下面的示例中,我們使用 KafkaProducer 實作向 Kafka 發送訊息的功能,在示例程式中,首先將 KafkaProduce 使用的配置寫入
到 Properties 中,每項配置的具體含義在注釋中進行解釋,之后以此 Properties 物件為引數構造 KafkaProducer 物件,最后通過 send 方法完成發送,代碼中包含同步發送、異步發送兩種情況,

在這里插入圖片描述
從上面的代碼可以看出 Kafka 為用戶提供了非常簡潔方便的 API,在使用時,只需要如下兩步:

  • 初始化 KafkaProducer 實體
  • 呼叫 send 介面發送資料

本文主要是圍繞著初始化 KafkaProducer 實體與如何實作 send 介面發送資料而展開的,

三、KafkaProducer 實體化

了解了 KafkaProducer 的基本使用,然后我們來深入了解下方法核心邏輯:

public KafkaProducer(Properties properties) {
    this(Utils.propsToMap(properties), (Serializer)null, (Serializer)null, (ProducerMetadata)null, (KafkaClient)null, (ProducerInterceptors)null, Time.SYSTEM);
}

在這里插入圖片描述

四、訊息發送程序

用戶是直接使用 producer.send() 發送的資料,先看一下 send() 介面的實作

// 異步向一個 topic 發送資料
public Future<RecordMetadata> send(ProducerRecord<K, V> record) {
    return this.send(record, (Callback)null);
}

// 向 topic 異步地發送資料,當發送確認后喚起回呼函式
public Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) {
    ProducerRecord<K, V> interceptedRecord = this.interceptors.onSend(record);
    return this.doSend(interceptedRecord, callback);
}

資料發送的最終實作還是呼叫了 Producer 的 doSend() 介面,

4.1 攔截器

首先方法會先進入攔截器集合 ProducerInterceptors , onSend 方法是遍歷攔截器 onSend 方 法,攔截器的目的是將資料處理加工, Kafka 本身并沒有給出默認的攔截器的實作,如果需要使用攔截器功能,必須自己實作介面,

4.1.1 攔截器代碼

在這里插入圖片描述
4.1.2 攔截器核心邏輯
在這里插入圖片描述
ProducerInterceptor 介面包括三個方法:

  • onSend(ProducerRecord<K, V> var1):該方法封裝進 KafkaProducer.send 方法中,即它運行在用戶主執行緒中的, 確保在訊息被序列化以計算磁區前呼叫該方法,用戶可以在該方法中對訊息做任何操作,但最好保證不要修改訊息所屬的 topic 和磁區,否則會影響目標磁區的計算,
  • onAcknowledgement(RecordMetadata var1, Exception var2):該方法會在訊息被應答之前或訊息發送失敗時呼叫,并且通常都是在 producer 回呼邏輯觸發之前,onAcknowledgement 運行在 producer 的 IO 執行緒中,因此不要在該方法中放入很重的邏輯,否則會拖慢 producer 的訊息發送效率,
  • close():關閉 interceptor,主要用于執行一些資源清理作業,

攔截器可能被運行在多個執行緒中,因此在具體實作時用戶需要自行確保執行緒安全,另外倘若指定了多個 interceptor,則 producer 將按照指定順序呼叫它們,并僅僅是捕獲每個 interceptor 可能拋出的例外記錄到錯誤日志中而非在向上傳遞,

4.2 Producer 的 doSend 實作

下面是 doSend() 的具體實作:

在這里插入圖片描述
在 doSend() 方法的實作上,一條 Record 資料的發送,主要分為以下五步:

  • 確認資料要發送到的 topic 的 metadata 是可用的(如果該 partition 的 leader 存在則是可用的,如果開啟權限時,client 有相應的權限),如果沒有 topic 的 metadata 資訊,就需要獲取相應的 metadata;
  • 序列化 record 的 key 和 value;
  • 獲取該 record 要發送到的 partition(可以指定,也可以根據演算法計算);
  • 向 accumulator 中追加 record 資料,資料會先進行快取;
  • 如果追加完資料后,對應的 RecordBatch 已經達到了 batch.size 的大小(或者 batch 的剩余空間不足以添加下一條 Record),則喚醒 sender 執行緒發送資料,

資料的發送程序,可以簡單總結為以上五點,下面會這幾部分的具體實作進行詳細分析,

五、訊息發送程序

5.1 獲取 topic 的 metadata 資訊

Producer 通過 waitOnMetadata() 方法來獲取對應 topic 的 metadata 資訊,這塊內容我下一篇再來講,

5.2 key 和 value 的序列化

Producer 端對 record 的 key 和 value 值進行序列化操作,在 Consumer 端再進行相應的反序列化,Kafka 內部提供的序列化和反序列化演算法如下圖所示:
在這里插入圖片描述
當然我們也是可以自定義序列化的具體實作,不過一般情況下,Kafka 內部提供的這些方法已經足夠使用,

5.3 獲取該 record 要發送到的 partition

獲取 partition 值,具體分為下面三種情況:

  • 指明 partition 的情況下,直接將指明的值直接作為 partiton 值;
  • 沒有指明 partition 值但有 key 的情況下,將 key 的 hash 值與 topic 的 partition 數進行取余得到 partition 值;
  • 既沒有 partition 值又沒有 key 值的情況下,第一次呼叫時隨機生成一個整數(后面每次呼叫在這個整數上自增),將這個值與 topic 可用的 partition 總數取余得到 partition 值,也就是常說的 round-robin 演算法,

具體實作如下:

// 當 record 中有 partition 值時,直接回傳,沒有的情況下呼叫 partitioner 的類的 partition 方法去計算(KafkaProducer.class)
private int partition(ProducerRecord<K, V> record, byte[] serializedKey, byte[] serializedValue, Cluster cluster) {
    Integer partition = record.partition();
    return partition != null ? partition : this.partitioner.partition(record.topic(), record.key(), serializedKey, record.value(), serializedValue, cluster);
}

Producer 默認使用的 partitioner 是 org.apache.kafka.clients.producer.internals.DefaultPartitioner,用戶也可以自定義 partition 的策略,下面是默認磁區策略具體實作:

public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
    return this.partition(topic, key, keyBytes, value, valueBytes, cluster, cluster.partitionsForTopic(topic).size());
}

public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster, int numPartitions) {
    return keyBytes == null ? this.stickyPartitionCache.partition(topic, cluster) : Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
}

在這里插入圖片描述
上面這個默認演算法核心就是粘著磁區快取

5.4 向 RecordAccmulator 中追加 record 資料

我們講 RecordAccumulator 之前先看這張圖,這樣的話會對整個發送流程有個大局觀,

在這里插入圖片描述
RecordAccmulator 承擔了緩沖區的角色,默認是 32 MB,

在 Kafka Producer 中,訊息不是一條一條發給 broker 的,而是多條訊息組成一個 ProducerBatch,然后由 Sender 一次性發出去,這里的 batch.size 并不是訊息的條數(湊滿多少條即發送),而是一個大小,默認是 16 KB,可以根據具體情況來進行優化,

在 RecordAccumulator 中,最核心的引數就是:

private final ConcurrentMap<TopicPartition, Deque<ProducerBatch>> batches;

它是一個 ConcurrentMap,key 是 TopicPartition 類,代表一個 topic 的一個 partition,value 是一個包含 ProducerBatch 的雙端佇列,等待 Sender 執行緒發送給 broker,畫張圖來看下:
在這里插入圖片描述

在這里插入圖片描述
上面的代碼不知道大家有沒有疑問?分配記憶體的代碼為啥不在 synchronized 同步塊中分配?導致下面的 synchronized 同步塊中還要 tryAppend 一下,

因為這時候可能其他執行緒已經創建好 RecordBatch 了,造成多余的記憶體申請,

如果把分配記憶體放在 synchronized 同步塊會有什么問題?

記憶體申請不到執行緒會一直等待,如果放在同步塊中會造成一直不釋放 Deque 佇列的鎖,那其他執行緒將無法對 Deque 佇列進行執行緒安全的同步操作,

再跟下 tryAppend() 方法,這就比較簡單了,

在這里插入圖片描述
以上代碼見圖解:

在這里插入圖片描述
5.5 喚醒 sender 執行緒發送 RecordBatch

當 record 寫入成功后,如果發現 RecordBatch 已滿足發送的條件(通常是 queue 中有多個 batch,那么最先添加的那些 batch 肯定是可以發送了),那么就會喚醒 sender 執行緒,發送 RecordBatch,

sender 執行緒對 RecordBatch 的處理是在 run() 方法中進行的,該方法具體實作如下:
在這里插入圖片描述
在這里插入圖片描述

其中比較核心的方法是 run() 方法中的 org.apache.kafka.clients.producer.internals.Sender#sendProducerData

其中 pollTimeout 意思是最長阻塞到至少有一個通道在你注冊的事件就緒了,回傳 0 則表示走起發車了,

在這里插入圖片描述
我們繼續跟下:org.apache.kafka.clients.producer.internals.RecordAccumulator#ready
在這里插入圖片描述
最后再來看下里面這個方法 org.apache.kafka.clients.producer.internals.RecordAccumulator#drain,從accumulator 緩沖區獲取要發送的資料,最大一次性發 max.request.size 大小的資料,

在這里插入圖片描述
在這里插入圖片描述

六、總結

最后為了讓你對 Kafka Producer 有個宏觀的架構理解,請看下圖:

在這里插入圖片描述
簡要說明:

  • new KafkaProducer() 后創建一個后臺執行緒 KafkaThread (實際運行執行緒是 Sender,KafkaThread 是對 Sender 的封裝) 掃描 RecordAccumulator 中是否有訊息,
  • 呼叫 KafkaProducer.send() 發送訊息,實際是將訊息保存到 RecordAccumulator 中,實際上就是保存到一個 Map 中 (ConcurrentMap<TopicPartition, Deque>),這條訊息會被記錄到同一個記錄批次 (相同主題相同磁區算同一個批次) 里面,這個批次的所有訊息會被發送到相同的主題和磁區上,
  • 后臺的獨立執行緒掃描到 RecordAccumulator 中有訊息后,會將訊息發送到 Kafka 集群中 (不是一有訊息就發送,而是要看訊息是否 ready)
  • 如果發送成功 (訊息成功寫入 Kafka), 就回傳一個 RecordMetaData 物件,它包括了主題和磁區資訊,以及記錄在磁區里的偏移量,
  • 如果寫入失敗,就會回傳一個錯誤,生產者在收到錯誤之后會嘗試重新發送訊息 (如果允許的話,此時會將訊息在保存到 RecordAccumulator 中),幾次之后如果還是失敗就回傳錯誤訊息,

好了,本文對 Kafka Producer 原始碼進行了決議,下一篇文章將會詳細介紹 metadata 的內容以及在 Producer 端 metadata 的更新機制,敬請期待~

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/298899.html

標籤:其他

上一篇:聽說看了這份Java學習路線的同學,畢業都拿到了大廠offer

下一篇:你們想知道的一切,都在這里了。

標籤雲
其他(157675) Python(38076) JavaScript(25376) Java(17977) C(15215) 區塊鏈(8255) C#(7972) AI(7469) 爪哇(7425) MySQL(7132) html(6777) 基礎類(6313) sql(6102) 熊猫(6058) PHP(5869) 数组(5741) R(5409) Linux(5327) 反应(5209) 腳本語言(PerlPython)(5129) 非技術區(4971) Android(4554) 数据框(4311) css(4259) 节点.js(4032) C語言(3288) json(3245) 列表(3129) 扑(3119) C++語言(3117) 安卓(2998) 打字稿(2995) VBA(2789) Java相關(2746) 疑難問題(2699) 细绳(2522) 單片機工控(2479) iOS(2429) ASP.NET(2402) MongoDB(2323) 麻木的(2285) 正则表达式(2254) 字典(2211) 循环(2198) 迅速(2185) 擅长(2169) 镖(2155) 功能(1967) .NET技术(1958) Web開發(1951) python-3.x(1918) HtmlCss(1915) 弹簧靴(1913) C++(1909) xml(1889) PostgreSQL(1872) .NETCore(1853) 谷歌表格(1846) Unity3D(1843) for循环(1842)

熱門瀏覽
  • 網閘典型架構簡述

    網閘架構一般分為兩種:三主機的三系統架構網閘和雙主機的2+1架構網閘。 三主機架構分別為內端機、外端機和仲裁機。三機無論從軟體和硬體上均各自獨立。首先從硬體上來看,三機都用各自獨立的主板、記憶體及存盤設備。從軟體上來看,三機有各自獨立的作業系統。這樣能達到完全的三機獨立。對于“2+1”系統,“2”分為 ......

    uj5u.com 2020-09-10 02:00:44 more
  • 如何從xshell上傳檔案到centos linux虛擬機里

    如何從xshell上傳檔案到centos linux虛擬機里及:虛擬機CentOs下執行 yum -y install lrzsz命令,出現錯誤:鏡像無法找到軟體包 前言 一、安裝lrzsz步驟 二、上傳檔案 三、遇到的問題及解決方案 總結 前言 提示:其實很簡單,往虛擬機上安裝一個上傳檔案的工具 ......

    uj5u.com 2020-09-10 02:00:47 more
  • 一、SQLMAP入門

    一、SQLMAP入門 1、判斷是否存在注入 sqlmap.py -u 網址/id=1 id=1不可缺少。當注入點后面的引數大于兩個時。需要加雙引號, sqlmap.py -u "網址/id=1&uid=1" 2、判斷文本中的請求是否存在注入 從文本中加載http請求,SQLMAP可以從一個文本檔案中 ......

    uj5u.com 2020-09-10 02:00:50 more
  • Metasploit 簡單使用教程

    metasploit 簡單使用教程 浩先生, 2020-08-28 16:18:25 分類專欄: kail 網路安全 linux 文章標簽: linux資訊安全 編輯 著作權 metasploit 使用教程 前言 一、Metasploit是什么? 二、準備作業 三、具體步驟 前言 Msfconsole ......

    uj5u.com 2020-09-10 02:00:53 more
  • 游戲逆向之驅動層與用戶層通訊

    驅動層代碼: #pragma once #include <ntifs.h> #define add_code CTL_CODE(FILE_DEVICE_UNKNOWN,0x800,METHOD_BUFFERED,FILE_ANY_ACCESS) /* 更多游戲逆向視頻www.yxfzedu.com ......

    uj5u.com 2020-09-10 02:00:56 more
  • 北斗電力時鐘(北斗授時服務器)讓網路資料更精準

    北斗電力時鐘(北斗授時服務器)讓網路資料更精準 北斗電力時鐘(北斗授時服務器)讓網路資料更精準 京準電子科技官微——ahjzsz 近幾年,資訊技術的得了快速發展,互聯網在逐漸普及,其在人們生活和生產中都得到了廣泛應用,并且取得了不錯的應用效果。計算機網路資訊在電力系統中的應用,一方面使電力系統的運行 ......

    uj5u.com 2020-09-10 02:01:03 more
  • 【CTF】CTFHub 技能樹 彩蛋 writeup

    ?碎碎念 CTFHub:https://www.ctfhub.com/ 筆者入門CTF時時剛開始刷的是bugku的舊平臺,后來才有了CTFHub。 感覺不論是網頁UI設計,還是題目質量,賽事跟蹤,工具軟體都做得很不錯。 而且因為獨到的金幣制度的確讓人有一種想去刷題賺金幣的感覺。 個人還是非常喜歡這個 ......

    uj5u.com 2020-09-10 02:04:05 more
  • 02windows基礎操作

    我學到了一下幾點 Windows系統目錄結構與滲透的作用 常見Windows的服務詳解 Windows埠詳解 常用的Windows注冊表詳解 hacker DOS命令詳解(net user / type /md /rd/ dir /cd /net use copy、批處理 等) 利用dos命令制作 ......

    uj5u.com 2020-09-10 02:04:18 more
  • 03.Linux基礎操作

    我學到了以下幾點 01Linux系統介紹02系統安裝,密碼啊破解03Linux常用命令04LAMP 01LINUX windows: win03 8 12 16 19 配置不繁瑣 Linux:redhat,centos(紅帽社區版),Ubuntu server,suse unix:金融機構,證券,銀 ......

    uj5u.com 2020-09-10 02:04:30 more
  • 05HTML

    01HTML介紹 02頭部標簽講解03基礎標簽講解04表單標簽講解 HTML前段語言 js1.了解代碼2.根據代碼 懂得挖掘漏洞 (POST注入/XSS漏洞上傳)3.黑帽seo 白帽seo 客戶網站被黑帽植入劫持代碼如何處理4.熟悉html表單 <html><head><title>TDK標題,描述 ......

    uj5u.com 2020-09-10 02:04:36 more
最新发布
  • 2023年最新微信小程式抓包教程

    01 開門見山 隔一個月發一篇文章,不過分。 首先回顧一下《微信系結手機號資料庫被脫庫事件》,我也是第一時間得知了這個訊息,然后跟蹤了整件事情的經過。下面是這起事件的相關截圖以及近日流出的一萬條資料樣本: 個人認為這件事也沒什么,還不如關注一下之前45億快遞資料查詢渠道疑似在近日復活的訊息。 訊息是 ......

    uj5u.com 2023-04-20 08:48:24 more
  • web3 產品介紹:metamask 錢包 使用最多的瀏覽器插件錢包

    Metamask錢包是一種基于區塊鏈技術的數字貨幣錢包,它允許用戶在安全、便捷的環境下管理自己的加密資產。Metamask錢包是以太坊生態系統中最流行的錢包之一,它具有易于使用、安全性高和功能強大等優點。 本文將詳細介紹Metamask錢包的功能和使用方法。 一、 Metamask錢包的功能 數字資 ......

    uj5u.com 2023-04-20 08:47:46 more
  • vulnhub_Earth

    前言 靶機地址->>>vulnhub_Earth 攻擊機ip:192.168.20.121 靶機ip:192.168.20.122 參考文章 https://www.cnblogs.com/Jing-X/archive/2022/04/03/16097695.html https://www.cnb ......

    uj5u.com 2023-04-20 07:46:20 more
  • 從4k到42k,軟體測驗工程師的漲薪史,給我看哭了

    清明節一過,盲猜大家已經無心上班,在數著日子準備過五一,但一想到銀行卡里的余額……瞬間心情就不美麗了。最近,2023年高校畢業生就業調查顯示,本科畢業月平均起薪為5825元。調查一出,便有很多同學表示自己又被平均了。看著這一資料,不免讓人想到前不久中國青年報的一項調查:近六成大學生認為畢業10年內會 ......

    uj5u.com 2023-04-20 07:44:00 more
  • 最新版本 Stable Diffusion 開源 AI 繪畫工具之中文自動提詞篇

    🎈 標簽生成器 由于輸入正向提示詞 prompt 和反向提示詞 negative prompt 都是使用英文,所以對學習母語的我們非常不友好 使用網址:https://tinygeeker.github.io/p/ai-prompt-generator 這個網址是為了讓大家在使用 AI 繪畫的時候 ......

    uj5u.com 2023-04-20 07:43:36 more
  • 漫談前端自動化測驗演進之路及測驗工具分析

    隨著前端技術的不斷發展和應用程式的日益復雜,前端自動化測驗也在不斷演進。隨著 Web 應用程式變得越來越復雜,自動化測驗的需求也越來越高。如今,自動化測驗已經成為 Web 應用程式開發程序中不可或缺的一部分,它們可以幫助開發人員更快地發現和修復錯誤,提高應用程式的性能和可靠性。 ......

    uj5u.com 2023-04-20 07:43:16 more
  • CANN開發實踐:4個DVPP記憶體問題的典型案例解讀

    摘要:由于DVPP媒體資料處理功能對存放輸入、輸出資料的記憶體有更高的要求(例如,記憶體首地址128位元組對齊),因此需呼叫專用的記憶體申請介面,那么本期就分享幾個關于DVPP記憶體問題的典型案例,并給出原因分析及解決方法。 本文分享自華為云社區《FAQ_DVPP記憶體問題案例》,作者:昇騰CANN。 DVPP ......

    uj5u.com 2023-04-20 07:43:03 more
  • msf學習

    msf學習 以kali自帶的msf為例 一、msf核心模塊與功能 msf模塊都放在/usr/share/metasploit-framework/modules目錄下 1、auxiliary 輔助模塊,輔助滲透(埠掃描、登錄密碼爆破、漏洞驗證等) 2、encoders 編碼器模塊,主要包含各種編碼 ......

    uj5u.com 2023-04-20 07:42:59 more
  • Halcon軟體安裝與界面簡介

    1. 下載Halcon17版本到到本地 2. 雙擊安裝包后 3. 步驟如下 1.2 Halcon軟體安裝 界面分為四大塊 1. Halcon的五個助手 1) 影像采集助手:與相機連接,設定相機引數,采集影像 2) 標定助手:九點標定或是其它的標定,生成標定檔案及內參外參,可以將像素單位轉換為長度單位 ......

    uj5u.com 2023-04-20 07:42:17 more
  • 在MacOS下使用Unity3D開發游戲

    第一次發博客,先發一下我的游戲開發環境吧。 去年2月份買了一臺MacBookPro2021 M1pro(以下簡稱mbp),這一年來一直在用mbp開發游戲。我大致分享一下我的開發工具以及使用體驗。 1、Unity 官網鏈接: https://unity.cn/releases 我一般使用的Apple ......

    uj5u.com 2023-04-20 07:40:19 more