当前位置: 首页 > news >正文

威海市城乡建设局网站给小学生做家教的网站

威海市城乡建设局网站,给小学生做家教的网站,力天装饰工程有限公司,金属质感 网站Kafka在Java项目中的应用 Docker 安装Kafka 一.首先需要安装docker,可看这篇文章安装docker 二.拉取zookeeper和KafKa镜像 docker pull wurstmeister/zookeeperdocker pull wurstmeister/kafkaKafka组件需要向zookeeper进行注册,所以也需要安装zookeeper 三.启动zookeeper…

Kafka在Java项目中的应用

Docker 安装Kafka

一.首先需要安装docker,可看这篇文章安装docker

二.拉取zookeeper和KafKa镜像

docker pull wurstmeister/zookeeperdocker pull wurstmeister/kafka

Kafka组件需要向zookeeper进行注册,所以也需要安装zookeeper

三.启动zookeeper、kafka组件

docker run -d --name zookeeper -p 2181:2181 wurstmeister/zookeeperdocker run -d --name kafka --publish 9092:9092 --link zookeeper --env KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 --env KAFKA_ADVERTISED_HOST_NAME=localhost --env KAFKA_ADVERTISED_PORT=9092 wurstmeister/kafka

启动成功界面如下,status即为running(运行中)
在这里插入图片描述

四.创建Springboot项目

4.1 添加依赖

<dependencies><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId></dependency><dependency><groupId>org.springframework.kafka</groupId><artifactId>spring-kafka</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-test</artifactId><scope>test</scope><exclusions><exclusion><groupId>org.junit.vintage</groupId><artifactId>junit-vintage-engine</artifactId></exclusion></exclusions></dependency><dependency><groupId>org.springframework.kafka</groupId><artifactId>spring-kafka-test</artifactId><scope>test</scope></dependency></dependencies>

4.2 application.yml文件

server:port: 9090
spring:kafka:bootstrap-servers: localhost:9092consumer:# 配置消费者消息offset是否自动重置(消费者重连会能够接收最开始的消息)auto-offset-reset: earliestproducer:value-serializer: org.springframework.kafka.support.serializer.JsonSerializerretries: 3  #  重试次数
kafka:topic:my-topic: my-topicmy-topic2: my-topic2

4.3 创建实体类Book

public class Book {private Long id;private String name;public Book() {}public Book(Long id, String name) {this.id = id;this.name = name;}public Long getId() {return id;}public void setId(Long id) {this.id = id;}public String getName() {return name;}public void setName(String name) {this.name = name;}@Overridepublic String toString() {return "Book{" +"id=" + id +", name='" + name + '\'' +'}';}
}

4.4 配置KafKa信息

@Configuration
public class KafkaConfig {@Value("${kafka.topic.my-topic}")String myTopic;@Value("${kafka.topic.my-topic2}")String myTopic2;/*** JSON消息转换器*/@Beanpublic RecordMessageConverter jsonConverter() {return new StringJsonMessageConverter();}/*** 通过注入一个 NewTopic 类型的 Bean 来创建 topic,如果 topic 已存在,则会忽略。*/@Beanpublic NewTopic myTopic() {return new NewTopic(myTopic, 2, (short) 1);}@Beanpublic NewTopic myTopic2() {return new NewTopic(myTopic2, 1, (short) 1);}
}

4.5 controller代码

@RestController
@RequestMapping(value = "/book")
public class BookController {@Value("${kafka.topic.my-topic}")String myTopic;@Value("${kafka.topic.my-topic2}")String myTopic2;BookProducerService producer;private AtomicLong atomicLong = new AtomicLong();BookController(BookProducerService producer) {this.producer = producer;}@GetMapping("/send")public String sendMessageToKafkaTopic(@RequestParam("name") String name) {this.producer.sendMessage(myTopic, new Book(atomicLong.addAndGet(1), name));this.producer.sendMessage(myTopic2, new Book(atomicLong.addAndGet(1), name));return name+" : 消息已经发送!";}
}

4.6 book 的生成者业务

@Service
public class BookProducerService {private static final Logger logger = LoggerFactory.getLogger(BookProducerService.class);private final KafkaTemplate<String, Object> kafkaTemplate;//通过构造方法进行注入public BookProducerService(KafkaTemplate<String, Object> kafkaTemplate) {this.kafkaTemplate = kafkaTemplate;}public void sendMessage(String topic, Object o) {ListenableFuture<SendResult<String, Object>> future = kafkaTemplate.send(topic, o);future.addCallback(result -> logger.info("生产者成功发送消息到topic:{} partition:{}的消息",result.getRecordMetadata().topic(),result.getRecordMetadata().partition()),ex -> logger.error("生产者发送消失败,原因:{}", ex.getMessage()));}}

4.7 book的消费者业务

@Service
public class BookConsumerService {@Value("${kafka.topic.my-topic}")private String myTopic;@Value("${kafka.topic.my-topic2}")private String myTopic2;private final Logger logger = LoggerFactory.getLogger(BookProducerService.class);private final ObjectMapper objectMapper = new ObjectMapper();@KafkaListener(topics = {"${kafka.topic.my-topic}"}, groupId = "group1")public void consumeMessage(ConsumerRecord<String, String> bookConsumerRecord) {try {Book book = objectMapper.readValue(bookConsumerRecord.value(), Book.class);logger.info("消费者消费topic:{} partition:{}的消息 -> {}", bookConsumerRecord.topic(), bookConsumerRecord.partition(), book.toString());} catch (JsonProcessingException e) {e.printStackTrace();}}@KafkaListener(topics = {"${kafka.topic.my-topic2}"}, groupId = "group2")public void consumeMessage2(Book book,ConsumerRecord<String,String> bookConsumerRecord) throws JsonProcessingException {Book value = objectMapper.readValue(bookConsumerRecord.value(), Book.class);logger.info("消费者消费topic:{} partition:{}的消息 -> {}", bookConsumerRecord.topic(), bookConsumerRecord.partition(), value.toString());logger.info("消费者消费{}的消息 -> {}", myTopic2, book.toString());}
}

代码整体目录如下

在这里插入图片描述

4.8 启动成功界面

在这里插入图片描述

4.9 浏览器访问

在这里插入图片描述

4.10 控制台显示

在这里插入图片描述

至此.基于KafKa的Springboot项目简单应用已经完成,后续需要对Kafka进行更深的学习以及应用!

http://www.yayakq.cn/news/758926/

相关文章:

  • 我想投诉做软件的网站公司名称变更
  • 网站里的聊天怎么做做一个响应式网站价格
  • 深圳网站制作问怎么自己做电商
  • 查找北京国互网网站建设空间网站大全
  • 设计网站建wordpress调用文章上级栏目名字
  • 我市强化属地网站建设电商要怎么做起来
  • 南通市城乡建设局网站河南省教育厅官方网站师德建设
  • 网站建设的市场有多大wordpress 压缩图片
  • 网站上的咨询窗口是怎么做的源码屋整站源码
  • wordpress建m域名网站企业网站建设企业
  • 公司网站可以分两个域名做吗举报不良网站信息怎么做
  • 化妆品企业网站建设前段模板网站
  • 绍兴做网站比较专业的公司大数据营销的特点有哪些
  • 自己做网站平台网站开发背景
  • 收到网站代码后怎么做短视频搜索seo
  • 建设网站青岛重庆建设公司网站
  • 甘井子区城市建设管理局网站用php做美食网站有哪些
  • ftp怎么做网站的备份王也踏青
  • 松岗做网站价格空间建网站
  • 制作网站微信登陆入口教我做网站
  • 网站建设问答网站后台管理系统 英文
  • 云霄县建设局网站wordpress怎么用panel
  • 网站安全狗卸载卸载不掉做网站主页上主要放哪些内容
  • 石家庄网站开发建设建设银行网上银行网站
  • 地方网站全网营销全球贸易中心网
  • 企业品牌网站建设类型教育培训机构网站源码
  • 北京怎么建设网站威海市建设工程协会网站
  • 网站做语言切换深圳宝安企业网站建设
  • 广西柳州网站建设推荐网站导航固定
  • 网站建设贵美工做网站是怎么做