spring集成kafka

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}
点赞
收藏

评论区

加载中...

相关推荐

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(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

手写Java HashMap源码

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

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )