結(jié)合線程池實(shí)現(xiàn)apache?kafka消費(fèi)者組的誤區(qū)及解決方法
一個(gè)錯(cuò)誤:多線程使用單一消費(fèi)者
下圖顯現(xiàn)了一種錯(cuò)誤的使用KafkaConsumer的方法
- 創(chuàng)建多個(gè)線程用來消費(fèi)kafka數(shù)據(jù)
- 多線程使用同一個(gè)KafkaConsumer對象
- 在單線程中使用這個(gè)KafkaConsumer對象,完成數(shù)據(jù)拉取、處理、提交偏移量。

這種方式之所以錯(cuò)誤的原因是:KafkaConsumer是線程不安全的,可能出現(xiàn)把同一批數(shù)據(jù)既給線程A處理,也交給線程B處理重復(fù)消費(fèi)的問題。
一個(gè)誤區(qū):多線程就是消費(fèi)者組
下圖中體現(xiàn)的是一種正常的KafkaConsumer使用方式
- 使用一個(gè)KafkaConsumer拉取數(shù)據(jù)
- 拉取數(shù)據(jù)后將一個(gè)批次的數(shù)據(jù)交給一個(gè)線程去處理

這個(gè)處理方式不是錯(cuò)誤,但是他只是一個(gè)消費(fèi)者在消費(fèi)kafka消息隊(duì)列中的數(shù)據(jù),不是消費(fèi)者組的方式消費(fèi)數(shù)據(jù)。無法充分利用kafka分區(qū)提升消息處理的吞吐量。
常規(guī)正確做法:使用線程池實(shí)現(xiàn)消費(fèi)者組
下面的方法是常規(guī)的正確實(shí)現(xiàn)方式

- 因?yàn)镵afkaConsumer是線程不安全的,所以不能跨線程使用KafkaConsumer
- 每個(gè)線程持有一個(gè)KafkaConsumer對象
- 多個(gè)線程的實(shí)現(xiàn)可以使用線程池,線程池的線程數(shù)量等于消費(fèi)者組內(nèi)消費(fèi)者的數(shù)量
public class MyConsumerGroup {
public void groupConsumer(){
ExecutorService executorService = Executors.newFixedThreadPool(6);
for (int i = 0; i < 6; i++) {
MyConsumer myConsumer = new MyConsumer();
executorService.execute(myConsumer);
}
}
}MyConsumer方法需要實(shí)現(xiàn)Runnable接口,并在run方法中調(diào)用MyConsumer#pollData。MyConsumer的代碼參考本專欄的《消費(fèi)者Java實(shí)現(xiàn)》( 集成apache kafka-clients實(shí)現(xiàn)數(shù)據(jù)消費(fèi)者)
@Override
public void run() {
MyConsumer myConsumer = new MyConsumer();
myConsumer.pollData();
}到此這篇關(guān)于結(jié)合線程池實(shí)現(xiàn)apache kafka消費(fèi)者組的文章就介紹到這了,更多相關(guān)apache kafka消費(fèi)者組內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java語言實(shí)現(xiàn)數(shù)據(jù)結(jié)構(gòu)棧代碼詳解
這篇文章主要介紹了Java語言實(shí)現(xiàn)數(shù)據(jù)結(jié)構(gòu)棧代碼詳解,簡單介紹了棧的概念,然后分享了線性棧和鏈?zhǔn)綏5腏ava代碼,具有一定參考價(jià)值,需要的朋友可以了解下。2017-11-11
ireport數(shù)據(jù)表格報(bào)表的簡單使用
這篇文章給大家介紹了如何畫一個(gè)報(bào)表模板,這里介紹下畫表格需要用到的組件,文中通過圖文并茂的形式給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧2021-10-10
Java 實(shí)現(xiàn)分布式服務(wù)的調(diào)用鏈跟蹤
分布式服務(wù)中完成某一個(gè)業(yè)務(wù)動(dòng)作,需要服務(wù)之間的相互協(xié)作才能完成,在這一次動(dòng)作引起的多服務(wù)的聯(lián)動(dòng)我們需要用1個(gè)唯一標(biāo)識(shí)關(guān)聯(lián)起來,關(guān)聯(lián)起來就是調(diào)用鏈的跟蹤。本文介紹了Java 實(shí)現(xiàn)分布式服務(wù)的調(diào)用鏈跟蹤的步驟2021-06-06
IDEA解決src和resource下創(chuàng)建多級目錄的操作
這篇文章主要介紹了IDEA解決src和resource下創(chuàng)建多級目錄的操作,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧2021-02-02

