Apache Kafka |
您所在的位置:网站首页 › kafka一次消费多少数据,怎么配置 › Apache Kafka |
![]() kafka提供了一些参数可以用于设置在消费端,用于提高消费的速度。 参数设置https://kafka.apache.org/24/documentation.html#consumerconfigs 支持的属性 见源码 KafkaProperties#Consumer 代码语言:javascript复制spring.kafka.listener.type 默认Single![]() ![]() 重点关注 ![]() 注意入参参数变为了 List 代码语言:javascript复制 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 messageMocks){ logger.info("【ArtisanCosumerMockDiffConsumeGroup接受到消息][线程:{} 消息大小:{}]", Thread.currentThread().getName(), messageMocks.size()); messageMocks.forEach(messageMock -> System.out.println("ArtisanCosumerMockDiffConsumeGroup收到的消息:" + messageMock)); } }单元测试代码语言:javascript复制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()); @Autowired private ArtisanProducerMock artisanProducerMock; @Test public void testAsynSend() throws ExecutionException, InterruptedException { logger.info("开始发送"); for (int i = 0; i < 2; i++) { artisanProducerMock.sendMsgASync().addCallback(new ListenableFutureCallback() { @Override public void onFailure(Throwable throwable) { logger.info(" 发送异常{}]]", throwable); } @Override public void onSuccess(SendResult objectObjectSendResult) { logger.info("回调结果 Result = topic:[{}] , partition:[{}], offset:[{}]", objectObjectSendResult.getRecordMetadata().topic(), objectObjectSendResult.getRecordMetadata().partition(), objectObjectSendResult.getRecordMetadata().offset()); } }); // 发送2次 每次间隔5秒, 凑够我们配置的 linger: ms: 10000 TimeUnit.SECONDS.sleep(5); logger.info("发送一条结束..."); } // 阻塞等待,保证消费 new CountDownLatch(1).await(); } }异步发送2条消息,每次发送消息之间, sleep 5 秒,以便达到配置的 linger.ms 最大等待时长10秒。 测试结果代码语言:javascript复制2021-02-18 12:13:00.201 INFO 8252 --- [ main] c.a.s.produceTest.ProduceMockTest : 开始发送 2021-02-18 12:13:05.426 INFO 8252 --- [ main] c.a.s.produceTest.ProduceMockTest : 发送一条结束... 2021-02-18 12:13:10.429 INFO 8252 --- [ main] c.a.s.produceTest.ProduceMockTest : 发送一条结束... 2021-02-18 12:13:10.442 INFO 8252 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回调结果 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 : 回调结果 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'}从日志中可以看出,发送的 2条消息被 消费者批量消费了 咦 , 我们把Type改成默认值试试呢? ![]() 重新测试 观察日志 代码语言:javascript复制2021-02-18 12:17:59.598 INFO 7764 --- [ main] c.a.s.produceTest.ProduceMockTest : 开始发送 2021-02-18 12:18:04.776 INFO 7764 --- [ main] c.a.s.produceTest.ProduceMockTest : 发送一条结束... 2021-02-18 12:18:09.778 INFO 7764 --- [ main] c.a.s.produceTest.ProduceMockTest : 发送一条结束... 2021-02-18 12:18:09.781 INFO 7764 --- [ad | producer-1] c.a.s.produceTest.ProduceMockTest : 回调结果 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 : 回调结果 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 |
今日新闻 |
推荐新闻 |
CopyRight 2018-2019 办公设备维修网 版权所有 豫ICP备15022753号-3 |