前往小程序,Get更优阅读体验!
立即前往
首页
学习
活动
专区
工具
TVP
发布
社区首页 >专栏 >Apache Kafka-消费端_批量消费消息的核心参数及功能实现

Apache Kafka-消费端_批量消费消息的核心参数及功能实现

作者头像
小小工匠
发布2021-08-17 16:38:53
2.6K0
发布2021-08-17 16:38:53
举报
文章被收录于专栏:小工匠聊架构

在这里插入图片描述
在这里插入图片描述

概述

kafka提供了一些参数可以用于设置在消费端,用于提高消费的速度。


参数设置

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

支持的属性 见源码 KafkaProperties#Consumer

代码语言:javascript
复制
spring.kafka.listener.type   默认Single
在这里插入图片描述
在这里插入图片描述
代码语言:javascript
复制
spring.kafka.consumer.max-poll-records
spring.kafka.consumer.fetch-min-size
spring.kafka.consumer.fetch-max-wait
在这里插入图片描述
在这里插入图片描述

Code

POM依赖

代码语言:javascript
复制
	<dependencies>
		<dependency>
			<groupId>org.springframework.bootgroupId>
			<artifactId>spring-boot-starter-webartifactId>
		dependency>

		
		<dependency>
			<groupId>org.springframework.kafkagroupId>
			<artifactId>spring-kafkaartifactId>
		dependency>

		<dependency>
			<groupId>org.springframework.bootgroupId>
			<artifactId>spring-boot-starter-testartifactId>
			<scope>testscope>
		dependency>
		<dependency>
			<groupId>junitgroupId>
			<artifactId>junitartifactId>
			<scope>testscope>
		dependency>
	dependencies>

配置文件

代码语言:javascript
复制
 spring:
  # Kafka 配置项,对应 KafkaProperties 配置类
  kafka:
    bootstrap-servers: 192.168.126.140:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔
    # Kafka Producer 配置项
    producer:
      acks: 1 # 0-不应答。1-leader 应答。all-所有 leader 和 follower 应答。
      retries: 3 # 发送失败时,重试发送的次数
      key-serializer: org.apache.kafka.common.serialization.StringSerializer # 消息的 key 的序列化
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # 消息的 value 的序列化
      batch-size: 16384 # 每次批量发送消息的最大数量   单位 字节  默认 16K
      buffer-memory: 33554432 # 每次批量发送消息的最大内存 单位 字节  默认 32M
      properties:
        linger:
          ms: 10000 # 批处理延迟时间上限。[实际不会配这么长,这里用于测速]这里配置为 10 * 1000 ms 过后,不管是否消息数量是否到达 batch-size 或者消息大小到达 buffer-memory 后,都直接发送一次请求。
    # Kafka Consumer 配置项
    consumer:
      auto-offset-reset: earliest # 设置消费者分组最初的消费进度为 earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring:
          json:
            trusted:
              packages: com.artisan.springkafka.domain
      fetch-max-wait: 10000  # poll 一次拉取的阻塞的最大时长,单位:毫秒。这里指的是阻塞拉取需要满足至少 fetch-min-size 大小的消息
      fetch-min-size: 10   # poll 一次消息拉取的最小数据量,单位:字节
      max-poll-records: 100  # poll 一次消息拉取的最大数量
    # Kafka Consumer Listener 监听器配置
    listener:
      missing-topics-fatal: false # 消费监听接口监听的主题不存在时,默认会报错。所以通过设置为 false ,解决报错
      type: batch # 监听器类型,默认为SINGLE ,只监听单条消息。这里我们配置 BATCH ,监听多条消息,批量消费

logging:
  level:
    org:
      springframework:
        kafka: ERROR # spring-kafka
      apache:
        kafka: ERROR # kafka

重点关注

在这里插入图片描述
在这里插入图片描述

生产者

代码语言:javascript
复制
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 {


    @Autowired
    private KafkaTemplate<Object,Object> kafkaTemplate ;


    /**
     * 同步发送
     * @return
     * @throws ExecutionException
     * @throws InterruptedException
     */
    public SendResult sendMsgSync() throws ExecutionException, InterruptedException {
        // 模拟发送的消息
        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 {
        // 模拟发送的消息
        Integer id = new Random().nextInt(100);
        MessageMock messageMock = new MessageMock(id,"messageSendByAsync-" + id);
        // 异步发送消息
        ListenableFuture<SendResult<Object, Object>> result = kafkaTemplate.send(TOPIC.TOPIC, messageMock);
        return result ;

    }

}

消费者

代码语言: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 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));

    }

}

注意入参参数变为了 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<MessageMock> 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<SendResult<Object, Object>>() {
                @Override
                public void onFailure(Throwable throwable) {
                    logger.info(" 发送异常{}]]", throwable);

                }
                @Override
                public void onSuccess(SendResult<Object, Object> 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

本文参与 腾讯云自媒体同步曝光计划,分享自作者个人站点/博客。
原始发表:2021/02/18 ,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 作者个人站点/博客 前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体同步曝光计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 概述
  • 参数设置
  • Code
    • POM依赖
      • 配置文件
        • 生产者
          • 消费者
            • 单元测试
              • 测试结果
              • 源码地址
              领券
              问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档