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){}}