springboot+kafka集成

kafka概念

  Topic:消息根据Topic进行归类,可以理解为一个队里。
  Producer:消息生产者,就是向kafka broker发消息的客户端。
  Consumer:消息消费者,向kafka broker取消息的客户端。
  broker:每个kafka实例(server),一台kafka服务器就是一个broker,一个集群由多个broker组成,一个broker可以容纳多个topic。
  Zookeeper:依赖集群保存meta信息。

1.安装

启动kafka需要zookeeper,所以需要安装zookeeper和kafka,这里就不介绍安装过程了。

2.引入kafka的依赖

1<dependencies> 2 <dependency> 3 <groupId>org.springframework.boot</groupId> 4 <artifactId>spring-boot-starter</artifactId> 5 </dependency> 6 7 <dependency> 8 <groupId>org.springframework.boot</groupId> 9 <artifactId>spring-boot-starter-test</artifactId> 10 <scope>test</scope> 11 </dependency> 12 13 <dependency> 14 <groupId>org.springframework.kafka</groupId> 15 <artifactId>spring-kafka</artifactId> 16 </dependency> 17 18 <dependency> 19 <groupId>com.google.code.gson</groupId> 20 <artifactId>gson</artifactId> 21 </dependency> 22 </dependencies>

3.定义消息

1public class Message { 2 private Long id; //id 3 4 private String msg; //消息 5 6 private Date sendTime; //时间戳 7 8 public Long getId() { 9 return id; 10 } 11 12 public void setId(Long id) { 13 this.id = id; 14 } 15 16 public String getMsg() { 17 return msg; 18 } 19 20 public void setMsg(String msg) { 21 this.msg = msg; 22 } 23 24 public Date getSendTime() { 25 return sendTime; 26 } 27 28 public void setSendTime(Date sendTime) { 29 this.sendTime = sendTime; 30 } 31 32}

4.提供者,注意加注解

发送消息时会有一个topic

1@Component 2public class KafkaSender { 3 4 private final org.slf4j.Logger log = LoggerFactory.getLogger(getClass()); 5 6 @Autowired 7 private KafkaTemplate<String, String> kafkaTemplate; 8 9 private Gson gson = new GsonBuilder().create(); 10 11 //发送消息方法 12 public void send() { 13 Message message = new Message(); 14 message.setId(System.currentTimeMillis()); 15 message.setMsg(UUID.randomUUID().toString()); 16 message.setSendTime(new Date()); 17 log.info("+++++++++++++++++++++ message = {}", gson.toJson(message)); 18 kafkaTemplate.send("zhisheng", gson.toJson(message)); 19 } 20}

5.消费者,注意加注解

需要加一个kafka监听的注解,里面需要指定topic,表示消费哪一种消息。

1@Component 2public class KafkaReceiver { 3 private final org.slf4j.Logger log = LoggerFactory.getLogger(getClass()); 4 5 @KafkaListener(topics = {"zhisheng"}) 6 public void listen(ConsumerRecord<?, ?> record) { 7 Optional<?> kafkaMessage = Optional.ofNullable(record.value()); 8 if (kafkaMessage.isPresent()) { 9 10 Object message = kafkaMessage.get(); 11 12 log.info("----------------- record =" + record); 13 log.info("------------------ message =" + message); 14 } 15 16 } 17}

7.启动类编写

这里调用消息提供方,写入了三条消息

1@SpringBootApplication 2public class KafkaDemoApplication { 3 4 public static void main(String[] args) { 5 ConfigurableApplicationContext context = SpringApplication.run(KafkaDemoApplication.class, args); 6 7 KafkaSender sender = context.getBean(KafkaSender.class); 8 9 for (int i = 0; i < 3; i++) { 10 //调用消息发送类中的消息发送方法 11 sender.send(); 12 13 try { 14 Thread.sleep(3000); 15 } catch (InterruptedException e) { 16 e.printStackTrace(); 17 } 18 } 19 } 20}

8.配置文件application.properties

1#============== kafka =================== 2# 指定kafka 代理地址,可以多个 3spring.kafka.bootstrap-servers=127.0.0.1:9092 4 5#=============== provider ======================= 6 7spring.kafka.producer.retries=0 8# 每次批量发送消息的数量 9spring.kafka.producer.batch-size=16384 10spring.kafka.producer.buffer-memory=33554432 11 12# 指定消息key和消息体的编解码方式 13spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer 14spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer 15 16#=============== consumer ======================= 17# 指定默认消费者group id 18spring.kafka.consumer.group-id=test-consumer-group 19 20spring.kafka.consumer.auto-offset-reset=earliest 21spring.kafka.consumer.enable-auto-commit=true 22spring.kafka.consumer.auto-commit-interval=100 23 24# 指定消息key和消息体的编解码方式 25spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer 26spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
点赞
收藏

评论区

加载中...

相关推荐

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(

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

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

作为一名程序员我不忘初心,复习指南

01kafka入门1.1什么是kafka 1.2kafka中的基本概念  1.2.1消息和批次  1.2.2主题和分区  1.2.3生产者和消费者、偏移量、消费者群组  1.2.4Broker和集群  1.2.5保留消息02为什么选择kafka2.1优点 2.2常见场景  2.2.1活动跟踪  2.2.2传递

Kafka安装步骤

基本概念1.Producer:消息生产者,就是向kafkabroker发消息的客户端2.Consumer:消息消费者,向kafkabroker取消息的客户端3.ConsumerGroup(CG):消费者组,由多个consumer组成。消费者组内每个消费者负责消费不同分区的数据,一个分区只能由一