聊聊如何利用kafka实现请求-响应模式

news/2025/2/8 2:55:58/文章来源:https://www.cnblogs.com/linyb-geek/p/18290423

前言

在大多数场景中,我们经常使用kafka来做发布-订阅,在发布-订阅模型中,消息一旦发送就不再追踪后续处理,但在某些业务场景下,我们希望在发送消息后等待一个响应,然后根据这个响应来做我们后续的操作。在这种请求-响应模式,我们就可以利用spring kafka的ReplyingKafkaTemplate来实现

ReplyingKafkaTemplate

简介

ReplyingKafkaTemplate 是 Spring Kafka 中的一个高级特性,专门用于处理 Kafka 中的请求/响应模式。它允许你发送一个消息到 Kafka,并等待一个响应

使用场景

  • 微服务间异步请求-响应: 当一个微服务需要从另一个微服务获取数据或执行操作,并希望在操作完成后得到通知时,可以使用 ReplyingKafkaTemplate。
  • 状态查询 如果一个服务需要定期或按需查询另一个服务的状态,但又不希望阻塞主线程等待响应,可以使用此模板。
  • 异步任务确认: 当一个服务发起一个异步任务(如文件上传、计算任务等),并需要知道任务何时完成时,可以使用 ReplyingKafkaTemplate 来接收完成通知

如何使用

1、在项目中引入spring-kafka gav

 <dependency><groupId>org.springframework.kafka</groupId><artifactId>spring-kafka</artifactId></dependency>

2、 配置replyingKafkaTemplate bean

注: ReplyingKafkaTemplate需依赖ProducerFactory和KafkaMessageListenerContainer

配置示例如下

 /*** 创建一个repliesContainer* @param containerFactory* @return*/@Beanpublic ConcurrentMessageListenerContainer<String, String> repliesContainer(ConcurrentKafkaListenerContainerFactory<String, String> containerFactory) {// 和RecordHeader中的topic对应起来ConcurrentMessageListenerContainer<String, String> repliesContainer =containerFactory.createContainer(KafkaConstant.REPLY_TOPIC);repliesContainer.getContainerProperties().setGroupId("repliesGroup");repliesContainer.setAutoStartup(false);return repliesContainer;}/*** 创建一个replyingTemplate* @param pf* @param repliesContainer* @return*/@Beanpublic ReplyingKafkaTemplate<String, String, String> replyingTemplate(ProducerFactory<String, String> pf,ConcurrentMessageListenerContainer<String, String> repliesContainer) {ReplyingKafkaTemplate<String, String, String> replyingKafkaTemplate = new ReplyingKafkaTemplate<>(pf, repliesContainer);// 设置响应超时为10秒,默认5秒replyingKafkaTemplate.setReplyTimeout(10000);return replyingKafkaTemplate;}

3、producer发送请求并等待响应

  @SneakyThrows@Overridepublic String sendAndReceive(String topic, ParamRequest request)  {// 创建ProducerRecord类,用来发送消息ProducerRecord<String,String> producerRecord = new ProducerRecord<>(topic, JSONUtil.toJsonStr(request));// 添加KafkaHeaders.REPLY_TOPIC到record的headers参数中,这个参数配置我们想要转发到哪个Topic中producerRecord.headers().add(new RecordHeader(KafkaHeaders.REPLY_TOPIC, KafkaConstant.REPLY_TOPIC.getBytes()));// sendAndReceive方法返回一个Future类RequestReplyFuture,// 这里类里面包含了获取发送结果的Future类和获取返回结果的Future类。// 使用replyingKafkaTemplate发送及返回都是异步操作RequestReplyFuture<String, String, String> replyFuture = replyingKafkaTemplate.sendAndReceive(producerRecord);// 获取发送结果SendResult<String, String> sendResult = replyFuture.getSendFuture().get();log.info("send message success,topic:{},message:{},sendResult:{}",topic,JSONUtil.toJsonStr(request),sendResult.getRecordMetadata());// 获取响应结果ConsumerRecord<String, String> consumerRecord = replyFuture.get();String result = consumerRecord.value();log.info("result: {}",result);return result;}

注: 方法里都写了相应注释,就不再论述了

4、consumer进行监听,并将返回结果通过@SendTo转发回去

在官网贴了这么一段话

他的大意是为了支持@SendTo,侦听器容器工厂必须提供一个KafkaTemplate(在其replyTemplate属性中),用于发送回复。因此我们做如下配置

   @Bean@ConditionalOnMissingBeanpublic ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(KafkaTemplate kafkaTemplate, ConsumerFactory consumerFactory) {ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();factory.setConsumerFactory(consumerFactory);factory.setReplyTemplate(kafkaTemplate);return factory;}

上述配置好后,就配置下监听器

 @KafkaListener(topics = TOPIC, groupId = "${lybgeek.consumer.group-id-prefix:lybgeek}-group-id", containerFactory = "kafkaListenerContainerFactory")/*** @SendTo 是一个Spring注解,常用于 Kafka 消费者方法之上,指示消息处理完成后应当将响应发送到哪个 Kafka 主题。* 使用场景:当你的应用作为服务端,需要对某个主题上的消息做出响应时,可以在处理该消息的方法上使用此注解来指定响应消息的目标主题。* 特点:简化了响应消息的路由配置,使得开发者无需显式地编写消息发送逻辑,只需关注业务处理逻辑。* 配合 ReplyingKafkaTemplate:在请求/响应模式中,@SendTo 指定的响应主题与 ReplyingKafkaTemplate 发送请求时设置的期望响应主题相匹配,从而使得请求方能够正确地接收响应消息。*///@SendTo("hello-test")@SendTopublic String listen(String data, Acknowledgment ack) {log.info("receive data:{}",data);if(JSONUtil.isJson(data)){Object result = execute(JSONUtil.toBean(data, ParamRequest.class),ack);if(result != null){ack.acknowledge();return JSONUtil.toJsonStr(result);};}return null;}

@SendTo的用途看我代码注释,具体用法可以查看官网
https://docs.spring.io/spring-kafka/reference/kafka/receiving-messages/annotation-send-to.html
进行了解

5、写个测试控制器

这个控制器的作用就是客户端发起http请求后,将请求参数送往kafka,kafka的消费方接收到http请求后,进行业务处理,并将业务结果通过kafka转发回去

  @PostMapping(value = "/**", consumes = {MediaType.APPLICATION_JSON_VALUE})
public Mono<ResponseEntity<byte[]>> forward(ProxyExchange<byte[]> proxy, Object params, HttpMethodEnum httpMethodEnum){try {String path = proxy.path().replace("/kafka", "").trim();ParamRequest paramRequest = buildParamRequest(path,params,httpMethodEnum);String topicCode = StringUtils.hasText(topicThreadLocal.get()) ? topicThreadLocal.get() : DEFAULT_TOPIC;String topic = TOPIC.replace(DEFAULT_TOPIC_PATTERN,topicCode);log.info(">>>>>>>>>>>>>> topic:{},httpMethod:{}, path:{},params:{}",topic,httpMethodEnum,path,params);Object result = kafkaService.sendAndReceive(topic,paramRequest);if(result != null){return Mono.just(ResponseEntity.ok(result.toString().getBytes()));}} catch (Exception e) {log.error(">>>>>>>>>>>>>> httpMethod:{},forward --> e:{}",httpMethodEnum.toString(),e.getMessage());} finally {topicThreadLocal.remove();}return Mono.just(ResponseEntity.ok(new byte[0]));}

核心就是这句代码

kafkaService.sendAndReceive(topic,paramRequest);

详细示例,可以查看文末的demo链接

使用ReplyingKafkaTemplate遇到的问题

No pending reply

这个问题是因为我拷贝了消费端配置文件,它配置了手动提交,而 ReplyingKafkaTemplate 是发送请求的一方,通常不需要特别的手动确认机制,因为 ReplyingKafkaTemplate 会等待响应或超时,因此改成自动确认即可

具体配置如下

spring:kafka:consumer:enable-auto-commit: ${KAFKA_CONSUMER_ENABLE_AUTO_COMMIT:true}

或者直接将消费端配置去掉也可以

总结

本文介绍通过ReplyingKafkaTemplate来实现请求-响应模式,在实际使用中,考虑到网络延迟和处理时间,调用ReplyingKafkaTemplate#sendAndReceive 方法可能会阻塞一段时间,因此在高负载环境下可能需要增加超时设置或使用回调机制

demo链接

https://github.com/lyb-geek/springboot-learning/tree/master/springboot-kafka-forward

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.hqwc.cn/news/845984.html

如若内容造成侵权/违法违规/事实不符,请联系编程知识网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

举个例子讲解DTO负责干啥

dto 在Spring Boot的开发过程中,使用DTO(Data Transfer Object)层是一个很常见的做法。DTO层是在应用程序的业务逻辑层和数据访问层之间引入的一个中间层,用于在不同层之间传输数据。本文将介绍DTO层的基本语法和为什么在Spring Boot开发中需要使用DTO层,并提供实际案例代…

39. css_01

1. css的概念 CSS(Cascading Style Sheets,层叠样式表)是一种用于描述HTML文档的表现形式的样式语言。它被设计用于将网页的内容与表现形式分离,可以控制网页的外观和布局,包括间距、颜色、字体等视觉元素,而不需要直接修改HTML的结构。 2. 语法结构选择符 {样式属性: 样…

DevExpress-独立使用的控件介绍-02

XtraEditors 库提供了只能独立使用的控件,即这些控件只能依附于其他控件配合使用,不能单独使用。这些控件包括:几种类型的列表控件、数据导航控件、滚动条和一个按钮控件,这些控件都是继承于BaseStyleControl,因此支持所有Dev 控件共有的样式、外观与感觉、以及工具提示机制…

preo/creo出图比例永久设置为1:1解决方法

平时画好的PREO/CREO模型需要转工程图后再导出到CAD时发现尺寸发生改变,那怎么设置可以导出比例1:1呢?下面分享下我的设置。 1.如图1画好的模型尺寸长宽为100*100MM的。(方法一)2.如图2当工程图再转CAD导出时,尺寸自动变成3.9370*3.9370MM。到时还需手动更改尺寸这样影响…

hhdb数据库介绍(10-29)

管理 数据备份 从存储节点或灾备机房数据备份 选择灾备机房类型、从库(双主备库)存储节点类型进行备份,页面根据选择类型,对应给出提示信息。发起备份时,检测从存储节点状态是否符合备份条件。主从数据一致性检测如果机房类型选择灾备机房或者存储节点类型选择从库(双主备…

hhdb数据库介绍(10-30)

管理 数据恢复 当业务数据遭受损坏或丢失时,可使用数据恢复功能将已备份的数据重新还原到损坏或丢失的逻辑库中。 数据恢复时序图:发起恢复发起说明点击“管理->数据恢复->【发起恢复】”即可跳转到数据恢复页面恢复发起前,出于数据安全性考虑,若超过3小时没有数据备…

使用Nginx搭建流媒体服务器

目录什么是流媒体服务器Nginx如何实现流媒体服务器为Nginx安装nginx-http-flv-module概述流程操作步骤配置流媒体服务器使用OBS推流使用VLC拉RTMP使用flv.js拉流使用jls.js拉m3u8总结引用 什么是流媒体服务器 流媒体服务器(Streaming Media Server)是一种用于存储和传输音频、…

[Go] Go语言教程

Go-lang概述Go 是一个开源的编程语言,它能让构造简单、可靠且高效的软件变得容易。 Go 是从2007年末由Robert Griesemer, Rob Pike, Ken Thompson主持开发,后来还加入了Ian Lance Taylor, Russ Cox等人,并最终于2009年11月开源,在2012年早些时候发布了Go 1稳定版本。现在Go…

hhdb数据库介绍(10-28)

管理 管理菜单主要囊括对业务数据进行管理的功能,例如对数据的备份恢复或执行业务表的DDL语句等操作。 数据对象 数据对象功能可以帮助用户通过列表实时查看当前已存在的数据对象,了解业务数据的整体情况。提供了对数据对象的筛选、统计、关联、详情等信息。基础数据对象的统…

hhdb数据库介绍(10-26)

报表 数据节点吞吐量 数据节点吞吐量为计算节点发往存储节点的操作量统计,一般用SELECT、UPDATE、DELETE、INSERT、OTHER五种类型分类计算节点操作。 图形模式 数据节点吞吐量图形模式包含数据节点吞吐总量对比图、数据节点吞吐量变化趋势、集群吞吐类型对比图、逻辑库吞吐量对…

离散数学命题逻辑

离散数学命题逻辑语雀链接:https://www.yuque.com/g/wushi-ls7km/zyko8c/tfttq5zq0xyldfxn/collaborator/join?token=u0bJmfKd8DcgpA1k&source=doc_collaborator# 《离散数学命题逻辑》

值班空岗睡岗识别智慧矿山视频分析技术安防摄像机的红外(补光)技术阐述科普

在现代安防监控领域,红外线(IR)技术因其在夜间或光线不足环境中的卓越表现而变得愈发重要。本文将深入探讨红外线技术在安防监控中的应用,分析其工作原理、分类以及在不同场景下的实际应用,同时探讨红外技术在智能交通和智慧矿山等领域中面临的挑战和解决方案。通过这一讨…