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