Springboot集成kafka實(shí)踐
在項(xiàng)目中使用kafka的場景有很多,尤其是實(shí)時(shí)產(chǎn)生的數(shù)據(jù)流,例如:電商數(shù)據(jù)、電信數(shù)據(jù)、統(tǒng)計(jì)等,通過kafka可以結(jié)合flink進(jìn)行大數(shù)據(jù)分析。所以第一步就是要集成kafka。
springboot已經(jīng)將kafka集成到框架里了,只需要引用依賴就可以簡單使用。
一、引入依賴
<!-- spring-kafka -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<version>2.8.6</version>
</dependency>二、在application.yml中配置信息
spring:
kafka:
# 以逗號(hào)分隔的地址列表,用于建立與Kafka集群的初始連接(kafka 默認(rèn)的端口號(hào)為9092)
bootstrap-servers: 192.168.1.199:9092
producer:
# 發(fā)生錯(cuò)誤后,消息重發(fā)的次數(shù)。
retries: 0
#當(dāng)有多個(gè)消息需要被發(fā)送到同一個(gè)分區(qū)時(shí),生產(chǎn)者會(huì)把它們放在同一個(gè)批次里。該參數(shù)指定了一個(gè)批次可以使用的內(nèi)存大小,按照字節(jié)數(shù)計(jì)算。
batch-size: 16384
# 設(shè)置生產(chǎn)者內(nèi)存緩沖區(qū)的大小。
buffer-memory: 33554432
# 鍵的序列化方式
key-serializer: org.apache.kafka.common.serialization.StringSerializer
# 值的序列化方式
value-serializer: org.apache.kafka.common.serialization.StringSerializer
# acks=0 : 生產(chǎn)者在成功寫入消息之前不會(huì)等待任何來自服務(wù)器的響應(yīng)。
# acks=1 : 只要集群的首領(lǐng)節(jié)點(diǎn)收到消息,生產(chǎn)者就會(huì)收到一個(gè)來自服務(wù)器成功響應(yīng)。
# acks=all :只有當(dāng)所有參與復(fù)制的節(jié)點(diǎn)全部收到消息時(shí),生產(chǎn)者才會(huì)收到一個(gè)來自服務(wù)器的成功響應(yīng)。
acks: all
consumer:
# 自動(dòng)提交的時(shí)間間隔 在spring boot 2.X 版本中這里采用的是值的類型為Duration 需要符合特定的格式,如1S,1M,2H,5D
auto-commit-interval: 1S
# 該屬性指定了消費(fèi)者在讀取一個(gè)沒有偏移量的分區(qū)或者偏移量無效的情況下該作何處理:
# latest(默認(rèn)值)在偏移量無效的情況下,消費(fèi)者將從最新的記錄開始讀取數(shù)據(jù)(在消費(fèi)者啟動(dòng)之后生成的記錄)
# earliest :在偏移量無效的情況下,消費(fèi)者將從起始位置讀取分區(qū)的記錄
auto-offset-reset: earliest
# 是否自動(dòng)提交偏移量,默認(rèn)值是true,為了避免出現(xiàn)重復(fù)數(shù)據(jù)和數(shù)據(jù)丟失,可以把它設(shè)置為false,然后手動(dòng)提交偏移量
enable-auto-commit: true
# 鍵的反序列化方式
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 值的反序列化方式
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
listener:
# 在偵聽器容器中運(yùn)行的線程數(shù)。
concurrency: 4三、 生產(chǎn)者
@RestController
@RequestMapping("/producer")
public class producerController {
@Autowired
private CallBackService callBackService;
@Autowired
KafkaTemplate<String, String> kafkaTemplate;
/**
* 列表
*/
@RequestMapping("/list")
public void list() {
List<CallBack> callBacks = callBackService.listMaps(10); // 從數(shù)據(jù)庫讀取10條記錄測試
for (CallBack callBack : callBacks) {
String json = JSONObject.toJSONString(callBack);
kafkaTemplate.send("callcdr", json); // 發(fā)送數(shù)據(jù)到kafka
System.out.println(json);
}
}
}四、消費(fèi)者
/**
* @Classname KafkaSimpleConsumer
* @Description 單個(gè)消費(fèi)者
* @Date 2019-05-14 10:08
*/
@Slf4j
@Component
public class KafkaConsumerTest {
// 簡單消費(fèi)者,groupId可以任意起
@KafkaListener(groupId = "ConsumerX", topics = "callcdr")
public void consumer1_1(ConsumerRecord<String, Object> record, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic, Consumer consumer) {
System.out.println("消費(fèi)記錄:" + record.value());
}
}五、使用BigData tools查看kafka情況
在IDEA里安裝BigData tools


消費(fèi)完成后consumerX為空。
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
springboot vue完成編輯頁面發(fā)送接口請(qǐng)求功能
這篇文章主要為大家介紹了springboot+vue完成編輯頁發(fā)送接口請(qǐng)求功能,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2022-05-05
SpringSecurity身份驗(yàn)證實(shí)現(xiàn)完整指南
本文介紹了如何使用SpringSecurity實(shí)現(xiàn)基于表單的登錄認(rèn)證,包括配置SecurityConfig、自定義登錄頁面、處理認(rèn)證成功和失敗的場景、獲取當(dāng)前登錄用戶、實(shí)現(xiàn)登出功能,幫助開發(fā)者理解和應(yīng)用SpringSecurity的核心功能和最佳實(shí)踐,感興趣的朋友跟隨小編一起看看吧2025-12-12
Java Swing JProgressBar進(jìn)度條的實(shí)現(xiàn)示例
這篇文章主要介紹了Java Swing JProgressBar進(jìn)度條的實(shí)現(xiàn)示例,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-12-12
Java程序啟動(dòng)時(shí)初始化數(shù)據(jù)的四種方式
本文主要介紹了Java程序啟動(dòng)時(shí)初始化數(shù)據(jù)的四種方式,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2024-02-02

