SpringBoot整合Kafka,通过简单配置实现生产消费功能
admin
2024-03-21 03:47:06
0

文章目录

  • 前提条件
  • 项目环境
  • 创建Topic
  • 配置信息
  • 生产消息
    • 生产自定义分区策略
    • 生产到指定分区
  • 消费消息
    • offset设置方式

*本文基于SpringBoot整合Kafka,通过简单配置实现生产及消费,包括生产消费的配置说明、消费者偏移设置方式等。更多功能细节可参考

spring kafka 文档:https://docs.spring.io/spring-kafka/docs/current/reference/html

前提条件

  • 搭建Kafka环境,参考Kafka集群环境搭建及使用
  • Java环境:JDK1.8
  • Maven版本:apache-maven-3.6.3
  • 开发工具:IntelliJ IDEA

项目环境

  1. 创建Springboot项目。
  2. pom.xml文件中引入kafka依赖。
org.springframework.kafkaspring-kafka

创建Topic

创建topic命名为testtopic并指定2个分区。

./kafka-topics.sh --bootstrap-server 127.0.0.1:9092 --create --topic testtopic --partitions 2

配置信息

application.yml配置文件信息

spring:application:name: kafka_springbootkafka:bootstrap-servers: 127.0.0.1:9092producer:#ACK机制,默认为1 (0,1,-1)acks: -1key-serializer: org.apache.kafka.common.serialization.StringSerializervalue-serializer: org.apache.kafka.common.serialization.StringSerializerproperties:# 自定义分区策略partitioner:class: org.bg.kafka.PartitionPolicyconsumer:#设置是否自动提交,默认为trueenable-auto-commit: falsekey-deserializer: org.apache.kafka.common.serialization.StringDeserializervalue-deserializer: org.apache.kafka.common.serialization.StringDeserializer#当一个新的消费组或者消费信息丢失后,在哪里开始进行消费。earliest:消费最早的消息。latest(默认):消费最近可用的消息。none:没有找到消费组消费数据时报异常。auto-offset-reset: latest#批量消费时每次poll的数量#max-poll-records: 5listener:#      当每一条记录被消费者监听器处理之后提交#      RECORD,#      当每一批数据被消费者监听器处理之后提交#      BATCH,#      当每一批数据被消费者监听器处理之后,距离上次提交时间大于TIME时提交#      TIME,#      当每一批数据被消费者监听器处理之后,被处理record数量大于等于COUNT时提交#      COUNT,#      #TIME | COUNT 有一个条件满足时提交#      COUNT_TIME,#      #当每一批数据被消费者监听器处理之后,手动调用Acknowledgment.acknowledge()后提交:#      MANUAL,#      # 手动调用Acknowledgment.acknowledge()后立即提交#      MANUAL_IMMEDIATE;ack-mode: manual#批量消费type: batch

更多配置信息查看KafkaProperties

生产消息

@Component
public class Producer {@Autowiredprivate KafkaTemplate kafkaTemplate;public void send(String msg) {kafkaTemplate.send(new ProducerRecord("testtopic", "key111", msg));}
}

生产自定义分区策略

package org.bg.kafka;import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.utils.Utils;import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.atomic.AtomicInteger;public class PartitionPolicy implements Partitioner {private final ConcurrentMap topicCounterMap = new ConcurrentHashMap();@Overridepublic int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {List partitions = cluster.partitionsForTopic(topic);int numPartitions = partitions.size();if (keyBytes == null) {int nextValue = this.nextValue(topic);List availablePartitions = cluster.availablePartitionsForTopic(topic);if (availablePartitions.size() > 0) {int part = Utils.toPositive(nextValue) % availablePartitions.size();return ((PartitionInfo)availablePartitions.get(part)).partition();} else {return Utils.toPositive(nextValue) % numPartitions;}} else {return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;}}private int nextValue(String topic) {AtomicInteger counter = (AtomicInteger)this.topicCounterMap.get(topic);if (null == counter) {counter = new AtomicInteger(ThreadLocalRandom.current().nextInt());AtomicInteger currentCounter = (AtomicInteger)this.topicCounterMap.putIfAbsent(topic, counter);if (currentCounter != null) {counter = currentCounter;}}return counter.getAndIncrement();}@Overridepublic void close() {}@Overridepublic void configure(Map map) {}
}

生产到指定分区

ProducerRecord有指定分区的构造方法,设置分区号
public ProducerRecord(String topic, Integer partition, K key, V value)

kafkaTemplate.send(new ProducerRecord("testtopic",1, "key111", msg));

消费消息


/*** 自定义seek参考* https://docs.spring.io/spring-kafka/docs/current/reference/html/#seek*/
@Component
public class Consumer implements ConsumerSeekAware{@KafkaListener(topics = {"testtopic"},groupId = "test_group",clientIdPrefix = "bg",id = "testconsumer")public void onMessage(List> records, Acknowledgment ack){System.out.println(records.size());System.out.println(records.toString());ack.acknowledge();}@Overridepublic void onPartitionsAssigned(Map assignments, ConsumerSeekCallback callback) {//按照时间戳设置偏移callback.seekToTimestamp(assignments.keySet(),1670233826705L);//设置偏移到最近callback.seekToEnd(assignments.keySet());//设置偏移到最开始callback.seekToBeginning(assignments.keySet());//指定 offsetfor (TopicPartition topicPartition : assignments.keySet()) {callback.seek(topicPartition.topic(),topicPartition.partition(),0L);}}}

offset设置方式

如代码所示,实现ConsumerSeekAware接口,设置offset几种方式:

  • 指定 offset,需要自己维护 offset,方便重试。
  • 指定从头开始消费。
  • 指定 offset 为最近可用的 offset (默认)。
  • 根据时间戳获取 offset,设置 offset。

相关内容

热门资讯

飞天茅台,又涨了100块,陈华... 作者:王一行 不到四个月,飞天茅台又涨价了。 加上3月31日那次涨价,今年飞天茅台的出厂价和零售价累...
银行理财收益缩水,机构集体喊话... 【大河财立方 记者 吴海舒 杨萨】“我自己买股票都没它能亏”,某社交平台上,一位网友晒出了自己购买的...
蒙商银行行长牛冠荣拟任内蒙古自... 蒙商银行行长牛冠荣拟任内蒙古自治区党委管理领导班子企业正职 人民财讯7月25日电,内蒙古自治区党委组...
首发经济破局 激活消费新动能 在昆明顺城购物中心,占地1800平方米的蜜雪冰城旗舰店人气爆棚,门口排满了前来打卡的消费者;蜡笔小新...
原创 谁... 坐在深圳南山的写字楼里往窗外看,无人机送外卖、机器人巡逻、满街的新能源车,很多外地人第一次来都会愣一...
“硬件创新基础设施”嘉立创今日... 7月24日,深圳嘉立创科技集团股份有限公司(以下简称“嘉立创”)正式启动网上网下发行申购,申购简称为...
深化产教融合 推进数智育人 哈... 7月23日,由阿里国际人工智能人才孵化中心(以下简称“阿里国际AITIC”)主办的“智启未来·数智赋...
陈春玉够“稳”,但魔法原子还“... 今年上半年,魔法原子获得了春晚的热度,但是也受到了人事和商业化的质疑。面对外界疑问,陈春玉依靠扎实的...
实物黄金和纸黄金的交易成本如何... 在黄金投资领域,实物黄金和纸黄金是较为常见的两种投资方式,而了解它们的交易成本计算方法对于投资者来说...
这些绩优股发布拟增持计划(附股... 7月以来,上市公司密集发布拟增持计划。与此同时,德明利、广钢气体、柯力传感等多家公司还发布了承诺不减...
特斯拉一周跌没18%,马斯克自... 马斯克这周不好过——特斯拉周五跌超2%,本周累跌近18%,创2022年以来最大单周跌幅;SpaceX...
原创 世... 文|江月白 编辑|江月白 近期中东局势再度掀起波澜,也门胡塞武装突然宣布封锁红海的曼德海峡,这一举...
农业农村部:乡村消费韧性持续凸... 本报记者 刘萌 7月24日,国新办举行新闻发布会介绍2026年上半年农业农村经济运行情况。农业农村部...
一杯鲜啤引爆夏夜狂欢 如东啤酒... 扬子晚报讯(记者 郭小川 通讯员 王军)如火的夏夜,怎能少了一杯清凉爽口的鲜啤?连日来,夜色中的如东...
原创 通... 时间定了,下周油价大涨!2026年汽柴油第10次上涨在即,时间将于7月31日24时准时调价,倒计时仅...
Waymo计划独立进入两地Ro... 7 月 25 日消息,据《金融时报》报道,Alphabet 旗下自动驾驶出租车企业 Waymo 在一...
日均狂赚2.39亿!宁德时代拿... 图片来源:图虫 7月24日晚,宁德时代(300750.SZ)披露2026年半年报,报告期内,公司实现...
长鑫科技下周一上市:合肥产投集... 长鑫科技下周一上市,大股东 合肥产投 都有哪些布局? 长鑫上市,合肥产投能赚多少? 作为长鑫科技发起...
ETF市场周报 | 市场回升趋... 市场回顾: 本周(2026年7月20日-7月24日),A股市场触底反弹,前4日整体走势强劲,周五略有...
美股开盘:指数涨跌不一 ,存储... 7月24日晚间,美股三大指数开盘后涨跌不一。截至发稿,标普500指数涨0.24%,道指涨0.32%,...