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

kafka生產(chǎn)者發(fā)送消息流程深入分析講解

 更新時(shí)間:2023年03月30日 09:17:06   作者:william_cr7  
本文將介紹kafka的一條消息的發(fā)送流程,從消息的發(fā)送到服務(wù)端的存儲(chǔ)。上文說到kafak分為客戶端與服務(wù)端,要發(fā)送消息就涉及到了網(wǎng)絡(luò)通訊,kafka采用TCP協(xié)議進(jìn)行客戶端與服務(wù)端的通訊協(xié)議

消息發(fā)送過程

消息的發(fā)送可能會(huì)經(jīng)過攔截器、序列化、分區(qū)器等過程。消息發(fā)送的主要涉及兩個(gè)線程,分別為main線程和sender線程。

如圖所示,主線程由 afkaProducer 創(chuàng)建消息,然后通過可能的攔截器、序列化器和分區(qū)器的作用之后緩存到消息累加器RecordAccumulator (也稱為消息收集器)中。 Sender 線程負(fù)責(zé)從RecordAccumulator 獲取消息并將其發(fā)送到 Kafka中。

攔截器

在消息序列化之前會(huì)經(jīng)過消息攔截器,自定義攔截器需要實(shí)現(xiàn)ProducerInterceptor接口,接口主要有兩個(gè)方案#onSend和#onAcknowledgement,在消息發(fā)送之前會(huì)調(diào)用前者方法,可以在發(fā)送之前假如處理邏輯,比如計(jì)費(fèi)。在收到服務(wù)端ack響應(yīng)后會(huì)觸發(fā)后者方法。需要注意的是攔截器中不要加入過多的復(fù)雜業(yè)務(wù)邏輯,以免影響發(fā)送效率。

消息分區(qū)

消息ProducerRecord會(huì)將消息路由到那個(gè)分區(qū)中,分兩種情況:

1.指定了partition字段

如果消息ProducerRecord中指定了 partition字段,那么就不需要走分區(qū)器,直接發(fā)往指定得partition分區(qū)中。

2.沒有指定partition,但自定義了分區(qū)器

3.沒指定parittion,也沒有自定義分區(qū)器,但key不為空

4.沒指定parittion,也沒有自定義分區(qū)器,key也為空

看源碼

// KafkaProducer#partition
private int partition(ProducerRecord<K, V> record, byte[] serializedKey, byte[] serializedValue, Cluster cluster) {
//指定分區(qū)partition則直接返回,否則走分區(qū)器
        Integer partition = record.partition();
        return partition != null ?
                partition :
                partitioner.partition(
                        record.topic(), record.key(), serializedKey, record.value(),                 serializedValue, cluster);
}
//DefaultPartitioner#partition
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        if (keyBytes == null) {
            return stickyPartitionCache.partition(topic, cluster);
        } 
        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();
        // hash the keyBytes to choose a partition
        return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
    }

partition 方法中定義了分區(qū)分配邏輯 如果 ke 不為 null , 那 么默認(rèn)的分區(qū)器會(huì)對(duì) key 進(jìn)行哈 希(采 MurmurHash2 算法 ,具備高運(yùn)算性能及 低碰 撞率),最終根據(jù)得到 哈希值來 算分區(qū)號(hào), 有相同 key 的消息會(huì)被寫入同一個(gè)分區(qū) 如果 key null ,那么消息將會(huì)以輪詢的方式發(fā)往主題內(nèi)的各個(gè)可用分區(qū)。

消息累加器

分區(qū)確定好了之后,消息并不是直接發(fā)送給broker,因?yàn)橐粋€(gè)個(gè)發(fā)送網(wǎng)絡(luò)消耗太大,而是先緩存到消息累加器RecordAccumulator,RecordAccumulator主要用來緩存消息 Sender 線程可以批量發(fā)送,進(jìn) 減少網(wǎng)絡(luò)傳輸 的資源消耗以提升性能 RecordAccumulator 緩存的大 小可以通過生產(chǎn)者客戶端參數(shù) buffer memory 配置,默認(rèn)值為 33554432B ,即 32MB如果生產(chǎn)者發(fā)送消息的速度超過發(fā) 送到服務(wù)器的速度 ,則會(huì)導(dǎo)致生產(chǎn)者空間不足,這個(gè)時(shí)候 KafkaProducer的send()方法調(diào)用要么 被阻塞,要么拋出異常,這個(gè)取決于參數(shù) max block ms 的配置,此參數(shù)的默認(rèn)值為 60秒。

消息累加器本質(zhì)上是個(gè)ConcurrentMap,

ConcurrentMap<TopicPartition, Deque<ProducerBatch>> batches;

發(fā)送流程源碼分析

//KafkaProducer
@Override
public Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) {
	// intercept the record, which can be potentially modified; this method does not throw exceptions
    //首先執(zhí)行攔截器鏈
	ProducerRecord<K, V> interceptedRecord = this.interceptors.onSend(record);
	return doSend(interceptedRecord, callback);
}
private Future<RecordMetadata> doSend(ProducerRecord<K, V> record, Callback callback) {
        TopicPartition tp = null;
	try {
		throwIfProducerClosed();
		// first make sure the metadata for the topic is available
		long nowMs = time.milliseconds();
		ClusterAndWaitTime clusterAndWaitTime;
		try {
			clusterAndWaitTime = waitOnMetadata(record.topic(), record.partition(), nowMs, maxBlockTimeMs);
		} catch (KafkaException e) {
			if (metadata.isClosed())
				throw new KafkaException("Producer closed while send in progress", e);
			throw e;
		}
		nowMs += clusterAndWaitTime.waitedOnMetadataMs;
		long remainingWaitMs = Math.max(0, maxBlockTimeMs - clusterAndWaitTime.waitedOnMetadataMs);
		Cluster cluster = clusterAndWaitTime.cluster;
		byte[] serializedKey;
		try {
			//key序列化
			serializedKey = keySerializer.serialize(record.topic(), record.headers(), record.key());
		} catch (ClassCastException cce) {
			throw new SerializationException("Can't convert key of class " + record.key().getClass().getName() +
					" to class " + producerConfig.getClass(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG).getName() +
					" specified in key.serializer", cce);
		}
		byte[] serializedValue;
		try {
			//value序列化
			serializedValue = valueSerializer.serialize(record.topic(), record.headers(), record.value());
		} catch (ClassCastException cce) {
			throw new SerializationException("Can't convert value of class " + record.value().getClass().getName() +
					" to class " + producerConfig.getClass(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG).getName() +
					" specified in value.serializer", cce);
		}
		//獲取分區(qū)partition
		int partition = partition(record, serializedKey, serializedValue, cluster);
		tp = new TopicPartition(record.topic(), partition);
		setReadOnly(record.headers());
		Header[] headers = record.headers().toArray();
		//消息壓縮
		int serializedSize = AbstractRecords.estimateSizeInBytesUpperBound(apiVersions.maxUsableProduceMagic(),
				compressionType, serializedKey, serializedValue, headers);
		//判斷消息是否超過最大允許大小,消息緩存空間是否已滿
		ensureValidRecordSize(serializedSize);
		long timestamp = record.timestamp() == null ? nowMs : record.timestamp();
		if (log.isTraceEnabled()) {
			log.trace("Attempting to append record {} with callback {} to topic {} partition {}", record, callback, record.topic(), partition);
		}
		// producer callback will make sure to call both 'callback' and interceptor callback
		Callback interceptCallback = new InterceptorCallback<>(callback, this.interceptors, tp);
 
		if (transactionManager != null && transactionManager.isTransactional()) {
			transactionManager.failIfNotReadyForSend();
		}
		//將消息緩存在消息累加器RecordAccumulator中
		RecordAccumulator.RecordAppendResult result = accumulator.append(tp, timestamp, serializedKey,
				serializedValue, headers, interceptCallback, remainingWaitMs, true, nowMs);
        //開辟新的ProducerBatch
		if (result.abortForNewBatch) {
			int prevPartition = partition;
			partitioner.onNewBatch(record.topic(), cluster, prevPartition);
			partition = partition(record, serializedKey, serializedValue, cluster);
			tp = new TopicPartition(record.topic(), partition);
			if (log.isTraceEnabled()) {
				log.trace("Retrying append due to new batch creation for topic {} partition {}. The old partition was {}", record.topic(), partition, prevPartition);
			}
			// producer callback will make sure to call both 'callback' and interceptor callback
			interceptCallback = new InterceptorCallback<>(callback, this.interceptors, tp);
 
			result = accumulator.append(tp, timestamp, serializedKey,
				serializedValue, headers, interceptCallback, remainingWaitMs, false, nowMs);
		}
		if (transactionManager != null && transactionManager.isTransactional())
			transactionManager.maybeAddPartitionToTransaction(tp);
		//判斷消息是否已滿,喚醒sender線程進(jìn)行發(fā)送消息
		if (result.batchIsFull || result.newBatchCreated) {
			log.trace("Waking up the sender since topic {} partition {} is either full or getting a new batch", record.topic(), partition);
			this.sender.wakeup();
		}
		return result.future;
		// handling exceptions and record the errors;
		// for API exceptions return them in the future,
		// for other exceptions throw directly
	} catch (Exception e) {
		// we notify interceptor about all exceptions, since onSend is called before anything else in this method
		this.interceptors.onSendError(record, tp, e);
		throw e;
	}
}

生產(chǎn)消息的可靠性

消息發(fā)送到broker,什么情況下生產(chǎn)者才確定消息寫入成功了呢?ack是生產(chǎn)者一個(gè)重要的參數(shù),它有三個(gè)值,ack=1表示leader副本寫入成功服務(wù)端即可返回給生產(chǎn)者,是吞吐量和消息可靠性的平衡方案;ack=0表示生產(chǎn)者發(fā)送消息之后不需要等服務(wù)端響應(yīng),這種消息丟失風(fēng)險(xiǎn)最大;ack=-1表示生產(chǎn)者需要等等ISR中所有副本寫入成功后才能收到響應(yīng),這種消息可靠性最高但吞吐量也是最小的。

到此這篇關(guān)于kafka生產(chǎn)者發(fā)送消息流程深入分析講解的文章就介紹到這了,更多相關(guān)kafka發(fā)送消息流程內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • springcloud+nacos實(shí)現(xiàn)灰度發(fā)布示例詳解

    springcloud+nacos實(shí)現(xiàn)灰度發(fā)布示例詳解

    這篇文章主要介紹了springcloud+nacos實(shí)現(xiàn)灰度發(fā)布,本文通過實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2023-08-08
  • Mac M1 Java 開發(fā)環(huán)境配置詳解

    Mac M1 Java 開發(fā)環(huán)境配置詳解

    這篇文章主要介紹了Mac M1 Java 開發(fā)環(huán)境配置詳解,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-03-03
  • Java字符串拼接的優(yōu)雅方式實(shí)例詳解

    Java字符串拼接的優(yōu)雅方式實(shí)例詳解

    字符串拼接一般使用“+”,但是“+”不能滿足大批量數(shù)據(jù)的處理,下面這篇文章主要給大家介紹了關(guān)于Java字符串拼接的幾種優(yōu)雅方式,需要的朋友可以參考下
    2021-07-07
  • springboot使用線程池(ThreadPoolTaskExecutor)示例

    springboot使用線程池(ThreadPoolTaskExecutor)示例

    大家好,本篇文章主要講的是springboot使用線程池(ThreadPoolTaskExecutor)示例,感興趣的同學(xué)趕快來看一看吧,對(duì)你有幫助的話記得收藏一下,方便下次瀏覽
    2021-12-12
  • logback的ShutdownHook關(guān)閉原理解析

    logback的ShutdownHook關(guān)閉原理解析

    這篇文章主要為大家介紹了logback的ShutdownHook關(guān)閉原理源碼解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-11-11
  • 淺談Spring Boot 開發(fā)REST接口最佳實(shí)踐

    淺談Spring Boot 開發(fā)REST接口最佳實(shí)踐

    這篇文章主要介紹了淺談Spring Boot 開發(fā)REST接口最佳實(shí)踐,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧
    2018-01-01
  • 詳解springboot項(xiàng)目docker部署實(shí)踐

    詳解springboot項(xiàng)目docker部署實(shí)踐

    這篇文章主要介紹了詳解springboot項(xiàng)目docker部署實(shí)踐,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧
    2018-01-01
  • 一篇文章帶你深入了解Java基礎(chǔ)(2)

    一篇文章帶你深入了解Java基礎(chǔ)(2)

    這篇文章主要給大家介紹了關(guān)于Java中方法使用的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2021-08-08
  • Java數(shù)據(jù)脫敏常用方法(3種)

    Java數(shù)據(jù)脫敏常用方法(3種)

    數(shù)據(jù)脫敏常用在電話號(hào)碼和身份證號(hào),本文主要介紹了Java數(shù)據(jù)脫敏常用方法,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2023-07-07
  • Spring Boot 3.0升級(jí)指南

    Spring Boot 3.0升級(jí)指南

    這篇文章主要為大家介紹了Spring Boot 3.0升級(jí)指南,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-02-02

最新評(píng)論

康保县| 沙湾县| 耒阳市| 新干县| 德昌县| 延津县| 定安县| 稷山县| 寿光市| 汉中市| 新和县| 汪清县| 富蕴县| 达日县| 上蔡县| 龙南县| 辽阳县| 鹤山市| 湘西| 宜兴市| 习水县| 东乡县| 天津市| 贵南县| 玉林市| 崇明县| 鹤山市| 红原县| 建湖县| 荣成市| 金山区| 增城市| 石嘴山市| 聂拉木县| 临高县| 准格尔旗| 江永县| 宁都县| 临湘市| 柏乡县| 万州区|