Springboot集成Kafka

 Kafka是一种高吞吐量的分布式发布订阅消息系统,有如下特性: 通过O(1)的磁盘数据结构提供消息的持久化,这种结构对于即使数以TB的消息存储也能够保持长时间的稳定性能。高吞吐量:即使是非常普通的硬件Kafka也可以支持每秒数百万的消息。支持通过Kafka服务器和消费机集群来分区消息。支持Hadoop并行数据加载。

Springboot的基本搭建和配置我在之前的文章已经给出代码示例了,如果还不了解的话可以先按照 SpringMVC配置太多?试试SpringBoot 进行学习哦。 那么如今很火的Springboot与kafka怎么完美的结合呢?多说无宜,放码过来 (talk is cheap,show me your code)!

安装Kafka

因为安装kafka需要zookeeper的支持,所以Windows安装时需要将zookeeper先安装上,然后将kafka安装好就可以了。 下面我给出Mac安装的步骤以及需要注意的点吧,windows的配置除了所在位置不太一样其他几乎没什么不同。

brew install kafka

对,就是这么简单,mac上一个命令就可以搞定了,这个安装过程可能需要等一会儿,应该是和网络状况有关系。安装提示信息可能有错误消息,如"Error: Could not link: /usr/local/share/doc/homebrew" 这个没关系,自动忽略掉了。 最终我们看到下面的样子就成功咯。

==> Summary 🍺/usr/local/Cellar/kafka/1.1.0: 157 files, 47.8MB

安装的配置文件位置如下,根据自己的需要修改端口号什么的就可以了。

安装的zoopeeper和kafka的位置 /usr/local/Cellar/

配置文件 /usr/local/etc/kafka/server.properties /usr/local/etc/kafka/zookeeper.properties

启动zookeeper

  ./bin/zookeeper-server-start /usr/local/etc/kafka/zookeeper.properties &

启动kafka 

./bin/kafka-server-start /usr/local/etc/kafka/server.properties &

为kafka创建Topic,topic 名为test,可以配置成自己想要的名字,回头再代码中配置正确就可以了。

 ./bin/kafka-topics --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic test

代码示例

pom.xml

1 <parent> 2 <groupId>org.springframework.boot</groupId> 3 <artifactId>spring-boot-starter-parent</artifactId> 4 <version>2.0.2.RELEASE</version> 5 <relativePath/> 6 </parent> 7 8 <properties> 9 <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> 10 <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> 11 <java.version>1.8</java.version> 12 </properties> 13 14 <dependencies> 15 <dependency> 16 <groupId>org.springframework.boot</groupId> 17 <artifactId>spring-boot-starter</artifactId> 18 </dependency> 19 <dependency> 20 <groupId>org.springframework.kafka</groupId> 21 <artifactId>spring-kafka</artifactId> 22 </dependency> 23 24 <dependency> 25 <groupId>org.springframework.boot</groupId> 26 <artifactId>spring-boot-starter-test</artifactId> 27 <scope>test</scope> 28 </dependency> 29 30 <dependency> 31 <groupId>com.google.code.gson</groupId> 32 <artifactId>gson</artifactId> 33 <version>2.8.2</version> 34 </dependency> 35 <dependency> 36 <groupId>org.springframework.boot</groupId> 37 <artifactId>spring-boot-starter-web</artifactId> 38 <version>RELEASE</version> 39 </dependency> 40 41 </dependencies>

application.yml

1server: 2 servlet: 3 context-path: / 4 port: 8080 5spring: 6 kafka: 7 bootstrap-servers: 127.0.0.1:9092 8 #生产者的配置,大部分我们可以使用默认的,这里列出几个比较重要的属性 9 producer: 10 #每批次发送消息的数量 11 batch-size: 16 12 #设置大于0的值将使客户端重新发送任何数据,一旦这些数据发送失败。注意,这些重试与客户端接收到发送错误时的重试没有什么不同。允许重试将潜在的改变数据的顺序,如果这两个消息记录都是发送到同一个partition,则第一个消息失败第二个发送成功,则第二条消息会比第一条消息出现要早。 13 retries: 0 14 #producer可以用来缓存数据的内存大小。如果数据产生速度大于向broker发送的速度,producer会阻塞或者抛出异常,以“block.on.buffer.full”来表明。这项设置将和producer能够使用的总内存相关,但并不是一个硬性的限制,因为不是producer使用的所有内存都是用于缓存。一些额外的内存会用于压缩(如果引入压缩机制),同样还有一些用于维护请求。 15 buffer-memory: 33554432 16 #key序列化方式 17 key-serializer: org.apache.kafka.common.serialization.StringSerializer 18 value-serializer: org.apache.kafka.common.serialization.StringSerializer 19 #消费者的配置 20 consumer: 21 #Kafka中没有初始偏移或如果当前偏移在服务器上不再存在时,默认区最新 ,有三个选项 【latest, earliest, none】 22 auto-offset-reset: latest 23 #是否开启自动提交 24 enable-auto-commit: true 25 #自动提交的时间间隔 26 auto-commit-interval: 100 27 #key的解码方式 28 key-deserializer: org.apache.kafka.common.serialization.StringDeserializer 29 #value的解码方式 30 value-deserializer: org.apache.kafka.common.serialization.StringDeserializer 31 #在/usr/local/etc/kafka/consumer.properties中有配置 32 group-id: test-consumer-group

Producer 消息生产者

1@Component 2public class Producer { 3 4 @Autowired 5 private KafkaTemplate kafkaTemplate; 6 7 private static Gson gson = new GsonBuilder().create(); 8 9 //发送消息方法 10 public void send() { 11 Message message = new Message(); 12 message.setId("KFK_"+System.currentTimeMillis()); 13 message.setMsg(UUID.randomUUID().toString()); 14 message.setSendTime(new Date()); 15 kafkaTemplate.send("test", gson.toJson(message)); 16 } 17 18} 19 20 21public class Message { 22 23 private String id; 24 25 private String msg; 26 27 private Date sendTime; 28 29 public String getId() { 30 return id; 31 } 32 33 public void setId(String id) { 34 this.id = id; 35 } 36 37 public String getMsg() { 38 return msg; 39 } 40 41 public void setMsg(String msg) { 42 this.msg = msg; 43 } 44 45 public Date getSendTime() { 46 return sendTime; 47 } 48 49 public void setSendTime(Date sendTime) { 50 this.sendTime = sendTime; 51 } 52} 53

Consumer 消息消费者

1public class Consumer { 2 3 @KafkaListener(topics = {"test"}) 4 public void listen(ConsumerRecord<?, ?> record){ 5 6 Optional<?> kafkaMessage = Optional.ofNullable(record.value()); 7 8 if (kafkaMessage.isPresent()) { 9 10 Object message = kafkaMessage.get(); 11 System.out.println("---->"+record); 12 System.out.println("---->"+message); 13 14 } 15 16 } 17} 18

测试接口用例

这里我们用一个接口来测试我们的消息发送会不会被消费者接收。

1@RestController 2@RequestMapping("/kafka") 3public class SendController { 4 5 @Autowired 6 private Producer producer; 7 8 @RequestMapping(value = "/send") 9 public String send() { 10 producer.send(); 11 return "{\"code\":0}"; 12 } 13} 14

在Springboot启动类启动后在浏览器访问http://127.0.0.1:8080/kafka/send,我们可以再IDE控制台中看到输出的结果,这时候我们的整合基本上就完成啦。 具体代码可以在SpringBootKafkaDemo@github获取哦。

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

Spring Boot 2.x 快速集成Kafka

1KafkaKafka是一个开源分布式的流处理平台,一种高吞吐量的分布式发布订阅消息系统,它可以处理消费者在网站中的所有动作流数据。Kafka由Scala和Java编写,2012年成为Apache基金会下顶级项目。2Kafka优点低延迟:Kafka支持低延迟消息传递,速度极快,能达到200w写/秒

FusionInsight大数据开发

Kafka应用开发1.了解Kafka应用开发适用场景2.熟悉Kafka应用开发流程3.熟悉并使用Kafka常用API4.进行Kafka应用开发Kafka的定义Kafka是一个高吞吐、分布式、基于发布订阅的消息系统Kafka有如下几个特点:1.高吞吐量2.消息持久化到磁

Kafka笔记

第1章Kafka简介1.1kafka起源Kafka是由LinkedIn开发并开源的分布式消息系统,2012年捐赠给Apache基金会,采用Scala语言,运行在JVM中,最新版本1.0.0。1.2kafka设计目标Kafka是一种分布式的,基于发布/订阅的消息系统。主要设计目标如下:①以时间复杂度O(1)的方式提供消息持久化能力,即