当前位置: 首页 > news >正文

Springboot集成kafka

Springboot集成kafka

环境配置

maven依赖如下 其余依赖根据业务自行配置

<dependency><groupId>org.springframework.kafka</groupId><artifactId>spring-kafka</artifactId><version>3.1.3</version><exclusions><exclusion><groupId>org.apache.kafka</groupId><artifactId>kafka-clients</artifactId></exclusion></exclusions></dependency><dependency><groupId>org.apache.kafka</groupId><artifactId>kafka-clients</artifactId><version>4.0.1</version><!-- 与你的 Kafka 服务端版本一致 --></dependency><dependency><groupId>org.jspecify</groupId><artifactId>jspecify</artifactId><version>1.0.0</version><scope>compile</scope></dependency><dependency><groupId>org.springframework</groupId><artifactId>spring-expression</artifactId><version>6.1.4</version><!-- 3.2.3对应的固定版本 --><scope>compile</scope><!-- 强制编译+运行时包含 --></dependency>

application.yml

spring:kafka:# bootstrap-servers: localhost:9092,localhost:9094,localhost:9096bootstrap-servers:localhost:9092# 生产者配置producer:# 消息Key/Value的序列化方式(必须和消费者反序列化对应)key-serializer:org.apache.kafka.common.serialization.StringSerializervalue-serializer:org.apache.kafka.common.serialization.StringSerializer# 消息确认机制:all=所有副本确认(最高可靠性,生产环境推荐)acks:all# 重试次数:发送失败时重试3次retries:3# 批次大小:16KB(批量发送提升性能)batch-size:16384# 缓冲区大小:32MBbuffer-memory:33554432# 消费者配置consumer:# 消息Key/Value的反序列化方式key-deserializer:org.apache.kafka.common.serialization.StringDeserializervalue-deserializer:org.apache.kafka.common.serialization.StringDeserializer# 消费者组ID(必填,同一组内的消费者负载均衡消费)group-id:springboot-kafka-group# 消费起始位置:earliest=从最开始消费;latest=从最新消息消费auto-offset-reset:earliest# 关闭自动提交偏移量(手动提交更安全,避免重复消费)enable-auto-commit:trueauto-commit-interval:1000

生产者producer

@ResourceprivateKafkaTemplate<String,String>kafkaTemplate;/** * 异步发送消息(带消息Key:相同Key的消息会进入同一个分区,保证有序) * @param topic 主题名 * @param key 消息Key * @param message 消息内容 */publicvoidsendMessageAsyncWithKey(Stringtopic,Stringkey,Stringmessage){ProducerRecord<String,String>record=newProducerRecord<>(topic,key,message);CompletableFuture<SendResult<String,String>>future=kafkaTemplate.send(record);future.thenAcceptAsync(newConsumer<SendResult<String,String>>(){@Overridepublicvoidaccept(SendResult<String,String>stringStringSendResult){System.out.println(stringStringSendResult.getProducerRecord());System.out.println(stringStringSendResult.getRecordMetadata());}});}

消费者consumer

@KafkaListener(topics="test-springboot-topic",groupId="springboot-kafka-group")publicvoidconsumeMessage(ConsumerRecord<String,String>record){System.out.println("消费到消息:"+record.value());}

创建topic

importorg.apache.kafka.clients.admin.Admin;importorg.apache.kafka.clients.admin.AdminClientConfig;importorg.apache.kafka.clients.admin.CreateTopicsResult;importorg.apache.kafka.clients.admin.NewTopic;importjava.util.Collections;importjava.util.Map;importjava.util.concurrent.ConcurrentHashMap;importjava.util.concurrent.ExecutionException;publicclassAdminTopicTest{publicstaticvoidmain(String[]args)throwsExecutionException,InterruptedException{Map<String,Object>configMap=newConcurrentHashMap<>();configMap.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,"127.0.0.1:9092,127.0.0.1:9094,127.0.0.1:9096");Adminadmin=Admin.create(configMap);//设置topic的名称 分区数 副本数NewTopictopic=newNewTopic("test111",3,(short)3);CreateTopicsResulttopics=admin.createTopics(Collections.singleton(topic));admin.close();}}

发送消息

import com.example.kafakademo.inceptor.FilterProducerInterceptor; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import java.util.HashMap; import java.util.Map; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.serialization.StringSerializer; public class KafkaproducerTest{public static void main(String[]args) throws ExecutionException,InterruptedException{Map<String,Object>configMap = new HashMap<>(); configMap.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"ip:9092"); configMap.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName()); configMap.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName()); configMap.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG,FilterProducerInterceptor.class.getName()); KafkaProducer<String,String>producer= new KafkaProducer<>(configMap); for (int i = 0; i < 10; i++){Future<RecordMetadata>send = producer.send(new ProducerRecord<>("test","key","xx"+i)); RecordMetadata recordMetadata = send.get(); System.out.println(recordMetadata.offset());}}}

消息拦截器

微服务链路追踪(traceId)
统一日志
消息加密
消息格式统一
全局异常处理
数据监控上报

import org.apache.kafka.clients.producer.ProducerInterceptor; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.springframework.util.StringUtils; import java.util.Map; public class FilterProducerInterceptor implements ProducerInterceptor<String,String>{@Override // 1. 消息发送之前调用(最常用) public ProducerRecord<String,String>onSend(ProducerRecord<String,String>record){if (StringUtils.isEmpty(record.value())){return null;}return new ProducerRecord<>(record.topic(),record.key(),"value:"+record.value());}@Override // 2. 消息发送成功/失败后调用 public void onAcknowledgement(RecordMetadata metadata,Exception exception){if (exception==null){System.out.println("success send topic "+metadata.topic()+" offset "+metadata.offset());}}@Override public void close(){}@Override public void configure(Map<String,?>configs){}}
http://www.cnnetsun.cn/news/1311890.html

相关文章:

  • 本科论文不用愁!Paperzz AI 毕业论文写作:四步搞定原创范文,图表公式全包含
  • 基于卷积神经网络-门控循环单元的时间序列预测 CNN-GRU 基于MATLAB环境 替换自己的...
  • 探秘风机螺栓预紧力的超声波检测:COMSOL 5.6仿真之旅
  • “南北合”随感
  • Vue+SpingBoot+MyBaits框架
  • 光电对抗:超构表面用于无源干扰的实践
  • AI简历分析:HR最在意的5大指标
  • 先甩个最核心的计数器代码镇楼
  • 【C++】C++类的幕后高手:友元、内部类、匿名对象与编译器优化深度解析
  • Python基于微信小程序的农产品溯源平台
  • 小白友好:Open-AutoGLM手机AI框架部署指南,10分钟跑通第一个自动化任务
  • Gemma-3-12b-it企业落地实践:中小团队低成本部署多模态AI助手方案
  • FR4 PCB透光LED反贴设计:丝印画中的隐藏式状态指示
  • ESP32-S3嵌入式HMI终端:LVGL驱动TFT+多模态感知设计
  • 支付宝周期扣款实战:从签约到代扣的全流程避坑指南(附代码示例)
  • TensorRT实战:FP16加速在边缘计算中的高效部署
  • STM32+Proteus仿真智能家居:手把手教你用虚拟串口模拟语音控制(附完整代码)
  • 构建网络安全培训系统:利用CosyVoice生成钓鱼邮件语音模拟与警示
  • Z-Image-Turbo_Sugar脸部Lora新手教程:零基础理解LoRA微调与Z-Image-Turbo底模关系
  • Mac用户必看:Chrome插件重启消失?3步永久解决同步问题(附图文)
  • 突破硬件限制:在PC上构建macOS开发环境的完整指南
  • gemma-3-12b-it惊艳效果展示:跨语言图文问答+多步推理真实案例集
  • PVE无线网络终极配置:让虚拟化主机摆脱网线束缚
  • Layuimini无限级菜单实战指南
  • ERNIE-4.5-0.3B-PT效果可视化:Chainlit中同一prompt不同温度值对比生成
  • 深入浅出GSM PDU 7bit编码:原理、实现与常见问题排查
  • Arduino结合PS2摇杆实现双舵机精准角度控制
  • 毫米波雷达拆解实录:从TI到大陆,硬件工程师带你读懂内部构造
  • Navicat for MySQL破解避坑指南:为什么你的PatchNavicat总是失败?
  • 独立开发者看过来:Z-Image-Turbo快速生成UI界面原型,节省外包成本