Apache Kafka-消费端_批量消费消息的核心参数及功能实现
生活随笔
收集整理的這篇文章主要介紹了
Apache Kafka-消费端_批量消费消息的核心参数及功能实现
小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,幫大家做個(gè)參考.
文章目錄
- 概述
- 參數(shù)設(shè)置
- Code
- POM依賴
- 配置文件
- 生產(chǎn)者
- 消費(fèi)者
- 單元測試
- 測試結(jié)果
- 源碼地址
概述
kafka提供了一些參數(shù)可以用于設(shè)置在消費(fèi)端,用于提高消費(fèi)的速度。
參數(shù)設(shè)置
https://kafka.apache.org/24/documentation.html#consumerconfigs
支持的屬性 見源碼 KafkaProperties#Consumer
spring.kafka.listener.type 默認(rèn)Single spring.kafka.consumer.max-poll-records spring.kafka.consumer.fetch-min-size spring.kafka.consumer.fetch-max-waitCode
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),對應(yīng) KafkaProperties 配置類kafka:bootstrap-servers: 192.168.126.140:9092 # 指定 Kafka Broker 地址,可以設(shè)置多個(gè),以逗號分隔# 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í)際不會配這么長,這里用于測速]這里配置為 10 * 1000 ms 過后,不管是否消息數(shù)量是否到達(dá) batch-size 或者消息大小到達(dá) buffer-memory 后,都直接發(fā)送一次請求。# 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.domainfetch-max-wait: 10000 # poll 一次拉取的阻塞的最大時(shí)長,單位:毫秒。這里指的是阻塞拉取需要滿足至少 fetch-min-size 大小的消息fetch-min-size: 10 # poll 一次消息拉取的最小數(shù)據(jù)量,單位:字節(jié)max-poll-records: 100 # poll 一次消息拉取的最大數(shù)量# Kafka Consumer Listener 監(jiān)聽器配置listener:missing-topics-fatal: false # 消費(fèi)監(jiān)聽接口監(jiān)聽的主題不存在時(shí),默認(rèn)會報(bào)錯(cuò)。所以通過設(shè)置為 false ,解決報(bào)錯(cuò)type: batch # 監(jiān)聽器類型,默認(rèn)為SINGLE ,只監(jiān)聽單條消息。這里我們配置 BATCH ,監(jiān)聽多條消息,批量消費(fèi)logging:level:org:springframework:kafka: ERROR # spring-kafkaapache:kafka: ERROR # kafka重點(diǎn)關(guān)注
生產(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;import java.util.List;/*** @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(List<MessageMock> messageMocks){logger.info("【ArtisanCosumerMock接受到消息][線程:{} 消息大小:{}]", Thread.currentThread().getName(), messageMocks.size());messageMocks.forEach(messageMock -> System.out.println("ArtisanCosumerMock收到的消息:" + messageMock));}}注意入?yún)?shù)變?yōu)榱?List
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;import java.util.List;/*** @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(List<MessageMock> messageMocks){logger.info("【ArtisanCosumerMockDiffConsumeGroup接受到消息][線程:{} 消息大小:{}]", Thread.currentThread().getName(), messageMocks.size());messageMocks.forEach(messageMock -> System.out.println("ArtisanCosumerMockDiffConsumeGroup收到的消息:" + messageMock));}}單元測試
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);logger.info("發(fā)送一條結(jié)束...");}// 阻塞等待,保證消費(fèi)new CountDownLatch(1).await();}}異步發(fā)送2條消息,每次發(fā)送消息之間, sleep 5 秒,以便達(dá)到配置的 linger.ms 最大等待時(shí)長10秒。
測試結(jié)果
2021-02-18 12:13:00.201 INFO 8252 --- [ main] c.a.s.produceTest.ProduceMockTest : 開始發(fā)送 2021-02-18 12:13:05.426 INFO 8252 --- [ main] c.a.s.produceTest.ProduceMockTest : 發(fā)送一條結(jié)束... 2021-02-18 12:13:10.429 INFO 8252 --- [ main] c.a.s.produceTest.ProduceMockTest : 發(fā)送一條結(jié)束... 2021-02-18 12:13:10.442 INFO 8252 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回調(diào)結(jié)果 Result = topic:[MOCK_TOPIC] , partition:[0], offset:[34] 2021-02-18 12:13:10.443 INFO 8252 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回調(diào)結(jié)果 Result = topic:[MOCK_TOPIC] , partition:[0], offset:[35] 2021-02-18 12:13:10.493 INFO 8252 --- [ntainer#0-0-C-1] c.a.s.consumer.ArtisanCosumerMock : 【ArtisanCosumerMock接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1 消息大小:2] 2021-02-18 12:13:10.493 INFO 8252 --- [ntainer#1-0-C-1] a.s.c.ArtisanCosumerMockDiffConsumeGroup : 【ArtisanCosumerMockDiffConsumeGroup接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#1-0-C-1 消息大小:2] ArtisanCosumerMockDiffConsumeGroup收到的消息:MessageMock{id=24, name='messageSendByAsync-24'} ArtisanCosumerMockDiffConsumeGroup收到的消息:MessageMock{id=32, name='messageSendByAsync-32'} ArtisanCosumerMock收到的消息:MessageMock{id=24, name='messageSendByAsync-24'} ArtisanCosumerMock收到的消息:MessageMock{id=32, name='messageSendByAsync-32'}從日志中可以看出,發(fā)送的 2條消息被 消費(fèi)者批量消費(fèi)了
咦 , 我們把Type改成默認(rèn)值試試呢?
重新測試
觀察日志
2021-02-18 12:17:59.598 INFO 7764 --- [ main] c.a.s.produceTest.ProduceMockTest : 開始發(fā)送 2021-02-18 12:18:04.776 INFO 7764 --- [ main] c.a.s.produceTest.ProduceMockTest : 發(fā)送一條結(jié)束... 2021-02-18 12:18:09.778 INFO 7764 --- [ main] c.a.s.produceTest.ProduceMockTest : 發(fā)送一條結(jié)束... 2021-02-18 12:18:09.781 INFO 7764 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回調(diào)結(jié)果 Result = topic:[MOCK_TOPIC] , partition:[0], offset:[36] 2021-02-18 12:18:09.782 INFO 7764 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回調(diào)結(jié)果 Result = topic:[MOCK_TOPIC] , partition:[0], offset:[37] 2021-02-18 12:18:09.837 INFO 7764 --- [ntainer#0-0-C-1] c.a.s.consumer.ArtisanCosumerMock : 【ArtisanCosumerMock接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1 消息大小:1] 2021-02-18 12:18:09.837 INFO 7764 --- [ntainer#1-0-C-1] a.s.c.ArtisanCosumerMockDiffConsumeGroup : 【ArtisanCosumerMockDiffConsumeGroup接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#1-0-C-1 消息大小:1] ArtisanCosumerMock收到的消息:MessageMock{id=13, name='messageSendByAsync-13'} ArtisanCosumerMockDiffConsumeGroup收到的消息:MessageMock{id=13, name='messageSendByAsync-13'} 2021-02-18 12:18:09.838 INFO 7764 --- [ntainer#0-0-C-1] c.a.s.consumer.ArtisanCosumerMock : 【ArtisanCosumerMock接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1 消息大小:1] ArtisanCosumerMock收到的消息:MessageMock{id=45, name='messageSendByAsync-45'} 2021-02-18 12:18:09.838 INFO 7764 --- [ntainer#1-0-C-1] a.s.c.ArtisanCosumerMockDiffConsumeGroup : 【ArtisanCosumerMockDiffConsumeGroup接受到消息][線程:org.springframework.kafka.KafkaListenerEndpointContainer#1-0-C-1 消息大小:1] ArtisanCosumerMockDiffConsumeGroup收到的消息:MessageMock{id=45, name='messageSendByAsync-45'}源碼地址
https://github.com/yangshangwei/boot2/tree/master/springkafkaBatchSend
總結(jié)
以上是生活随笔為你收集整理的Apache Kafka-消费端_批量消费消息的核心参数及功能实现的全部內(nèi)容,希望文章能夠幫你解決所遇到的問題。
- 上一篇: Apache Kafka-生产者_批量发
- 下一篇: Apache Kafka-通过设置Con