1、引入依赖jar包
1<dependency> 2 <groupId>org.springframework.kafka</groupId> 3 <artifactId>spring-kafka</artifactId> 4</dependency>
2、配置kafka信息
1spring: 2 kafka: 3 bootstrap-servers: localhost:9092 4 consumer: 5 group-id: group1 6 listener: 7 missing-topics-fatal: false
启动报错,需要配置missing-topics-fatal为false
1org.springframework.context.ApplicationContextException: Failed to start bean 'org.springframework.kafka.config.internalKafkaListenerEndpointRegistry'; nested exception is java.lang.IllegalStateException: Topic(s) [topic1] is/are not present and missingTopicsFatal is true 2 at org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:185) 3 at org.springframework.context.support.DefaultLifecycleProcessor.access$200(DefaultLifecycleProcessor.java:53) 4 at org.springframework.context.support.DefaultLifecycleProcessor$LifecycleGroup.start(DefaultLifecycleProcessor.java:360) 5 at org.springframework.context.support.DefaultLifecycleProcessor.startBeans(DefaultLifecycleProcessor.java:158) 6 at org.springframework.context.support.DefaultLifecycleProcessor.onRefresh(DefaultLifecycleProcessor.java:122) 7 at org.springframework.context.support.AbstractApplicationContext.finishRefresh(AbstractApplicationContext.java:894) 8 at org.springframework.boot.web.servlet.context.ServletWebServerApplicationContext.finishRefresh(ServletWebServerApplicationContext.java:162) 9 at org.springframework.context.support.AbstractApplicationContext.refresh(AbstractApplicationContext.java:553) 10 at org.springframework.boot.web.servlet.context.ServletWebServerApplicationContext.refresh(ServletWebServerApplicationContext.java:141) 11 at org.springframework.boot.SpringApplication.refresh(SpringApplication.java:747) 12 at org.springframework.boot.SpringApplication.refreshContext(SpringApplication.java:397) 13 at org.springframework.boot.SpringApplication.run(SpringApplication.java:315) 14 at org.springframework.boot.SpringApplication.run(SpringApplication.java:1226) 15 at org.springframework.boot.SpringApplication.run(SpringApplication.java:1215) 16 at com.zhoulp.SyncMessageWebApplication.main(SyncMessageWebApplication.java:24) 17Caused by: java.lang.IllegalStateException: Topic(s) [topic1] is/are not present and missingTopicsFatal is true 18 at org.springframework.kafka.listener.AbstractMessageListenerContainer.checkTopics(AbstractMessageListenerContainer.java:383) 19 at org.springframework.kafka.listener.ConcurrentMessageListenerContainer.doStart(ConcurrentMessageListenerContainer.java:144) 20 at org.springframework.kafka.listener.AbstractMessageListenerContainer.start(AbstractMessageListenerContainer.java:340) 21 at org.springframework.kafka.config.KafkaListenerEndpointRegistry.startIfNecessary(KafkaListenerEndpointRegistry.java:312) 22 at org.springframework.kafka.config.KafkaListenerEndpointRegistry.start(KafkaListenerEndpointRegistry.java:257) 23 at org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:182) 24 ... 14 common frames omitted
3、实现生产者
1package com.zhoulp.producer; 2 3import javax.inject.Inject; 4 5import org.slf4j.Logger; 6import org.slf4j.LoggerFactory; 7import org.springframework.kafka.core.KafkaTemplate; 8import org.springframework.stereotype.Component; 9 10/** 11 * 12 * @author zhoulp 13 * @date 2020-08-03 14 * 15 */ 16@Component("kafkaProducer") 17public class KafkaProducer { 18 19 private static Logger log = LoggerFactory.getLogger(KafkaProducer.class); 20 21 @Inject 22 private KafkaTemplate<String, String> template; 23 24 public void sendMessage(String topic, String data) { 25 log.info("send: topic = {}, data = {}", topic, data); 26 template.send(topic, data); 27 } 28 29}
4、实现消费者
1package com.zhoulp.consumer; 2 3import org.apache.kafka.clients.consumer.ConsumerRecord; 4import org.slf4j.Logger; 5import org.slf4j.LoggerFactory; 6import org.springframework.kafka.annotation.KafkaListener; 7import org.springframework.stereotype.Component; 8 9/** 10 * 11 * @author zhoulp 12 * @date 2020-08-03 13 * 14 */ 15@Component("kafkaConsumer") 16public class KafkaConsumer { 17 18 private static Logger log = LoggerFactory.getLogger(KafkaConsumer.class); 19 20 @KafkaListener(topics = "topic1") 21 public void listenTopic1(ConsumerRecord<String, String> consumerRecord) { 22 log.info("listenTopic1"); 23 log.info(consumerRecord.toString()); 24 log.info(consumerRecord.topic()); 25 log.info(consumerRecord.value()); 26 } 27 28}