日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當(dāng)前位置: 首頁 > 编程资源 > 编程问答 >内容正文

编程问答

Apache Kafka-生产者_批量发送消息的核心参数及功能实现

發(fā)布時(shí)間:2025/3/21 编程问答 25 豆豆
生活随笔 收集整理的這篇文章主要介紹了 Apache Kafka-生产者_批量发送消息的核心参数及功能实现 小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,幫大家做個(gè)參考.

文章目錄

  • 概述
  • 參數(shù)設(shè)置
  • Code
    • POM依賴
    • 配置文件
    • 生產(chǎn)者
    • 消費(fèi)者
    • 單元測(cè)試
    • 測(cè)試結(jié)果
  • 源碼地址


概述

kafka中有個(gè) micro batch 的概念 ,為了提高Producer 發(fā)送的性能。

不同于RocketMQ 提供了一個(gè)可以批量發(fā)送多條消息的 API 。 Kafka 的做法是:提供了一個(gè) RecordAccumulator 消息收集器,將發(fā)送給相同 Topic 的相同 Partition 分區(qū)的消息們,緩沖一下,當(dāng)滿足條件時(shí)候,一次性批量將緩沖的消息提交給 Kafka Broker 。


參數(shù)設(shè)置

https://kafka.apache.org/24/documentation.html#producerconfigs

主要涉及的參數(shù) ,三個(gè)條件,滿足任一即會(huì)批量發(fā)送:

  • batch-size :超過收集的消息數(shù)量的最大量。默認(rèn)16KB

  • buffer-memory :超過收集的消息占用的最大內(nèi)存 , 默認(rèn)32M

  • linger.ms :超過收集的時(shí)間的最大等待時(shí)長(zhǎng),單位:毫秒。


Code

POM依賴

<dependencies><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId></dependency><!-- 引入 Spring-Kafka 依賴 --><dependency><groupId>org.springframework.kafka</groupId><artifactId>spring-kafka</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-test</artifactId><scope>test</scope></dependency><dependency><groupId>junit</groupId><artifactId>junit</artifactId><scope>test</scope></dependency></dependencies>

配置文件

spring:# Kafka 配置項(xiàng),對(duì)應(yīng) KafkaProperties 配置類kafka:bootstrap-servers: 192.168.126.140:9092 # 指定 Kafka Broker 地址,可以設(shè)置多個(gè),以逗號(hào)分隔# Kafka Producer 配置項(xiàng)producer:acks: 1 # 0-不應(yīng)答。1-leader 應(yīng)答。all-所有 leader 和 follower 應(yīng)答。retries: 3 # 發(fā)送失敗時(shí),重試發(fā)送的次數(shù)key-serializer: org.apache.kafka.common.serialization.StringSerializer # 消息的 key 的序列化value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # 消息的 value 的序列化batch-size: 16384 # 每次批量發(fā)送消息的最大數(shù)量 單位 字節(jié) 默認(rèn) 16Kbuffer-memory: 33554432 # 每次批量發(fā)送消息的最大內(nèi)存 單位 字節(jié) 默認(rèn) 32Mproperties:linger:ms: 10000 # 批處理延遲時(shí)間上限。[實(shí)際不會(huì)配這么長(zhǎng),這里用于測(cè)速]這里配置為 10 * 1000 ms 過后,不管是否消息數(shù)量是否到達(dá) batch-size 或者消息大小到達(dá) buffer-memory 后,都直接發(fā)送一次請(qǐng)求。# Kafka Consumer 配置項(xiàng)consumer:auto-offset-reset: earliest # 設(shè)置消費(fèi)者分組最初的消費(fèi)進(jìn)度為 earliestkey-deserializer: org.apache.kafka.common.serialization.StringDeserializervalue-deserializer: org.springframework.kafka.support.serializer.JsonDeserializerproperties:spring:json:trusted:packages: com.artisan.springkafka.domain# Kafka Consumer Listener 監(jiān)聽器配置listener:missing-topics-fatal: false # 消費(fèi)監(jiān)聽接口監(jiān)聽的主題不存在時(shí),默認(rèn)會(huì)報(bào)錯(cuò)。所以通過設(shè)置為 false ,解決報(bào)錯(cuò)logging:level:org:springframework:kafka: ERROR # spring-kafkaapache:kafka: ERROR # kafka


生產(chǎn)者

package com.artisan.springkafka.producer;import com.artisan.springkafka.constants.TOPIC; import com.artisan.springkafka.domain.MessageMock; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; import org.springframework.stereotype.Component; import org.springframework.util.concurrent.ListenableFuture;import java.util.Random; import java.util.concurrent.ExecutionException;/*** @author 小工匠* @version 1.0* @description: TODO* @date 2021/2/17 22:25* @mark: show me the code , change the world*/@Component public class ArtisanProducerMock {@Autowiredprivate KafkaTemplate<Object,Object> kafkaTemplate ;/*** 同步發(fā)送* @return* @throws ExecutionException* @throws InterruptedException*/public SendResult sendMsgSync() throws ExecutionException, InterruptedException {// 模擬發(fā)送的消息Integer id = new Random().nextInt(100);MessageMock messageMock = new MessageMock(id,"artisanTestMessage-" + id);// 同步等待return kafkaTemplate.send(TOPIC.TOPIC, messageMock).get();}public ListenableFuture<SendResult<Object, Object>> sendMsgASync() throws ExecutionException, InterruptedException {// 模擬發(fā)送的消息Integer id = new Random().nextInt(100);MessageMock messageMock = new MessageMock(id,"messageSendByAsync-" + id);// 異步發(fā)送消息ListenableFuture<SendResult<Object, Object>> result = kafkaTemplate.send(TOPIC.TOPIC, messageMock);return result ;}}

消費(fèi)者

package com.artisan.springkafka.consumer;import com.artisan.springkafka.domain.MessageMock; import com.artisan.springkafka.constants.TOPIC; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component;/*** @author 小工匠* @version 1.0* @description: TODO* @date 2021/2/17 22:33* @mark: show me the code , change the world*/@Component public class ArtisanCosumerMock {private Logger logger = LoggerFactory.getLogger(getClass());private static final String CONSUMER_GROUP_PREFIX = "MOCK-A" ;@KafkaListener(topics = TOPIC.TOPIC ,groupId = CONSUMER_GROUP_PREFIX + TOPIC.TOPIC)public void onMessage(MessageMock messageMock){logger.info("【接受到消息][線程:{} 消息內(nèi)容:{}]", Thread.currentThread().getName(), messageMock);}} package com.artisan.springkafka.consumer;import com.artisan.springkafka.domain.MessageMock; import com.artisan.springkafka.constants.TOPIC; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component;/*** @author 小工匠* @version 1.0* @description: TODO* @date 2021/2/17 22:33* @mark: show me the code , change the world*/@Component public class ArtisanCosumerMockDiffConsumeGroup {private Logger logger = LoggerFactory.getLogger(getClass());private static final String CONSUMER_GROUP_PREFIX = "MOCK-B" ;@KafkaListener(topics = TOPIC.TOPIC ,groupId = CONSUMER_GROUP_PREFIX + TOPIC.TOPIC)public void onMessage(MessageMock messageMock){logger.info("【接受到消息][線程:{} 消息內(nèi)容:{}]", Thread.currentThread().getName(), messageMock);}}

單元測(cè)試

package com.artisan.springkafka.produceTest;import com.artisan.springkafka.SpringkafkaApplication; import com.artisan.springkafka.producer.ArtisanProducerMock; import org.junit.Test; import org.junit.runner.RunWith; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.kafka.support.SendResult; import org.springframework.test.context.junit4.SpringRunner; import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.ListenableFutureCallback;import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit;/*** @author 小工匠* * @version 1.0* @description: TODO* @date 2021/2/17 22:40* @mark: show me the code , change the world*/@RunWith(SpringRunner.class) @SpringBootTest(classes = SpringkafkaApplication.class) public class ProduceMockTest {private Logger logger = LoggerFactory.getLogger(getClass());@Autowiredprivate ArtisanProducerMock artisanProducerMock;@Testpublic void testAsynSend() throws ExecutionException, InterruptedException {logger.info("開始發(fā)送");for (int i = 0; i < 2; i++) {artisanProducerMock.sendMsgASync().addCallback(new ListenableFutureCallback<SendResult<Object, Object>>() {@Overridepublic void onFailure(Throwable throwable) {logger.info(" 發(fā)送異常{}]]", throwable);}@Overridepublic void onSuccess(SendResult<Object, Object> objectObjectSendResult) {logger.info("回調(diào)結(jié)果 Result = topic:[{}] , partition:[{}], offset:[{}]",objectObjectSendResult.getRecordMetadata().topic(),objectObjectSendResult.getRecordMetadata().partition(),objectObjectSendResult.getRecordMetadata().offset());}});// 發(fā)送2次 每次間隔5秒, 湊夠我們配置的 linger: ms: 10000TimeUnit.SECONDS.sleep(5);}// 阻塞等待,保證消費(fèi)new CountDownLatch(1).await();}}

異步發(fā)送2條消息,每次發(fā)送消息之間, sleep 5 秒,以便達(dá)到配置的 linger.ms 最大等待時(shí)長(zhǎng)10秒。


測(cè)試結(jié)果

2021-02-18 10:58:53.360 INFO 24736 --- [ main] c.a.s.produceTest.ProduceMockTest : 開始發(fā)送 2021-02-18 10:59:03.555 INFO 24736 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回調(diào)結(jié)果 Result = topic:[MOCK_TOPIC] , partition:[0], offset:[30] 2021-02-18 10:59:03.556 INFO 24736 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回調(diào)結(jié)果 Result = topic:[MOCK_TOPIC] , partition:[0], offset:[31] 2021-02-18 10:59:03.595 INFO 24736 --- [ntainer#0-0-C-1] c.a.s.consumer.ArtisanCosumerMock : 【接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1 消息內(nèi)容:MessageMock{id=6, name='messageSendByAsync-6'}] 2021-02-18 10:59:03.595 INFO 24736 --- [ntainer#1-0-C-1] a.s.c.ArtisanCosumerMockDiffConsumeGroup : 【接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#1-0-C-1 消息內(nèi)容:MessageMock{id=6, name='messageSendByAsync-6'}] 2021-02-18 10:59:03.595 INFO 24736 --- [ntainer#1-0-C-1] a.s.c.ArtisanCosumerMockDiffConsumeGroup : 【接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#1-0-C-1 消息內(nèi)容:MessageMock{id=94, name='messageSendByAsync-94'}] 2021-02-18 10:59:03.595 INFO 24736 --- [ntainer#0-0-C-1] c.a.s.consumer.ArtisanCosumerMock : 【接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1 消息內(nèi)容:MessageMock{id=94, name='messageSendByAsync-94'}]

10 秒后,滿足批量消息的最大等待時(shí)長(zhǎng),所以 2 條消息被 Producer 批量發(fā)送。同時(shí)我們配置的是 acks=1 ,需要等待發(fā)送成功后,才會(huì)回調(diào) ListenableFutureCallback 的方法。

當(dāng)然了,我們這里都是為了測(cè)試,設(shè)置的這么長(zhǎng)的間隔,實(shí)際中需要根據(jù)具體的業(yè)務(wù)場(chǎng)景設(shè)置一個(gè)合理的值。


源碼地址

https://github.com/yangshangwei/boot2/tree/master/springkafkaBatchSend

總結(jié)

以上是生活随笔為你收集整理的Apache Kafka-生产者_批量发送消息的核心参数及功能实现的全部?jī)?nèi)容,希望文章能夠幫你解決所遇到的問題。

如果覺得生活随笔網(wǎng)站內(nèi)容還不錯(cuò),歡迎將生活随笔推薦給好友。

主站蜘蛛池模板: 欧美另类色 | 精品无码av在线 | 日韩成人综合网 | 2021亚洲天堂 | 久久国产精品电影 | 色综合99久久久无码国产精品 | 6090伦理 | 天天躁日日躁狠狠躁av麻豆男男 | 午夜激情国产 | 夜夜夜久久久 | 国产日日干 | 亚洲综合久久av一区二区三区 | 久草欧美视频 | 亚洲第一区在线播放 | 91一区二区三区在线观看 | 黄a网站 | av激情影院 | 欧美成人黑人xx视频免费观看 | 99精品欧美一区二区蜜桃免费 | av网站免费播放 | jjzzjjzz欧美69巨大 | 亚洲中文字幕视频一区 | 精品人妻无码一区 | 男性影院 | 国产精品福利一区二区三区 | 理论片在线观看视频 | 游戏涩涩免费网站 | 欧美一级做a爰片久久高潮 久热国产精品视频 | 九九热视频在线 | 麻豆传媒一区二区三区 | 中国毛片在线 | 九色亚洲| 播播开心激情网 | 成人中文字幕在线 | 尤物网在线 | 麻豆影视网站 | 国产精品亚洲二区 | 国产精品s色 | 色综合久久88色综合天天6 | 亚洲美女自拍偷拍 | 狠狠爱亚洲 | 一级做a爰片 | 毛片毛片毛片毛片毛片毛片毛片毛片毛片毛片 | 动漫av一区二区 | 一区二区三区免费播放 | 国内一区二区视频 | 三级黄片毛片 | 亚洲图片自拍偷拍区 | 美女大逼| 日韩人妻精品一区二区三区 | 男人的天堂在线观看av | 老汉av网站 | 少妇av一区二区三区 | 午夜精品久久久久久久爽 | 欧洲一区二区在线观看 | 天堂精品视频 | 黄色片免费视频 | 青青插 | 一区二区日韩在线观看 | 国产精品综合一区二区 | 欧美日b片 | 亚洲激情综合 | 国产天堂视频 | 久久久久久视 | a级黄色录像 | 国产三级精品在线 | a视频在线观看免费 | av中文字幕网址 | 日本精品久久久久久 | 国产91在线播放 | 国产a级免费视频 | 一区二区欧美在线观看 | 久久91亚洲 | 超级黄色片| 色国产视频 | 国产福利观看 | 国产污污视频在线观看 | 91女神在线 | 91av在线免费 | 麻豆系列| 又大又粗弄得我出好多水 | 国产小视频在线免费观看 | ass亚洲肉体欣赏pics | 一级做a爱片久久 | 久久精品久久久久久久 | xxxxx在线| 国产欧美视频在线观看 | 欧美一级做a爰片免费视频 成人激情在线观看 | 婷婷爱五月| 97视频精品| 污片网址| 国产小精品 | 俺去日| 日韩中文一区二区三区 | 色狠狠干| 长腿校花无力呻吟娇喘的视频 | 强行糟蹋人妻hd中文字幕 | 欧美草逼网 | 无码少妇一区二区 |