基于ReplyingKafkaTemplate实现REST调用的请求-响应模式方案咨询
基于ReplyingKafkaTemplate实现REST调用的请求-响应模式方案咨询
嗨,刚好我对这个场景很熟悉,给你一套基于ReplyingKafkaTemplate的落地方案,完美适配你的REST同步调用+Kafka消息流转需求~
先明确下你的场景细节:
- 服务A:对外提供同步REST接口,接收外部客户端的请求
- 服务B:监听
RequestTopic处理请求,处理完成后把结果发送到ResponseTopic - 核心流程:外部客户端→(同步REST)服务A→(发送请求)
RequestTopic→服务B处理→(发送结果)ResponseTopic→服务A接收结果→(返回响应)外部客户端
具体实现步骤
1. 服务A侧配置与实现
ReplyingKafkaTemplate是Spring Kafka专门为请求-响应模式设计的模板,它能帮你自动关联请求和响应,实现同步等待结果的逻辑。
首先是配置类,需要定义生产者工厂、响应监听容器和ReplyingKafkaTemplate实例:
import org.springframework.kafka.core.*; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.support.serializer.JsonSerializer; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaConfig { // 配置Kafka生产者工厂 @Bean public ProducerFactory<String, YourRequestDto> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址:9092"); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); return new DefaultKafkaProducerFactory<>(configProps); } // 配置响应主题的监听容器,用于接收服务B返回的结果 @Bean public ConcurrentMessageListenerContainer<String, YourResponseDto> replyContainer(ConsumerFactory<String, YourResponseDto> consumerFactory) { ContainerProperties containerProps = new ContainerProperties("ResponseTopic"); return new ConcurrentMessageListenerContainer<>(consumerFactory, containerProps); } // 实例化ReplyingKafkaTemplate @Bean public ReplyingKafkaTemplate<String, YourRequestDto, YourResponseDto> replyingKafkaTemplate(ProducerFactory<String, YourRequestDto> producerFactory, ConcurrentMessageListenerContainer<String, YourResponseDto> replyContainer) { return new ReplyingKafkaTemplate<>(producerFactory, replyContainer); } }
然后在REST接口里使用这个模板,实现同步发送请求并等待响应:
import org.springframework.http.ResponseEntity; import org.springframework.kafka.requestreply.RequestReplyFuture; import org.springframework.kafka.core.ReplyingKafkaTemplate; import org.springframework.kafka.support.ProducerRecord; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RestController; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeoutException; import java.util.concurrent.TimeUnit; @RestController @RequestMapping("/external") public class ExternalApiController { private final ReplyingKafkaTemplate<String, YourRequestDto, YourResponseDto> replyingKafkaTemplate; // 构造注入模板 public ExternalApiController(ReplyingKafkaTemplate<String, YourRequestDto, YourResponseDto> replyingKafkaTemplate) { this.replyingKafkaTemplate = replyingKafkaTemplate; } @PostMapping("/process") public ResponseEntity<YourResponseDto> handleExternalRequest(@RequestBody YourRequestDto request) throws InterruptedException, ExecutionException, TimeoutException { // 构造请求消息,指定发送到RequestTopic ProducerRecord<String, YourRequestDto> requestRecord = new ProducerRecord<>("RequestTopic", request); // 发送请求并设置超时时间(比如5秒,根据业务调整) RequestReplyFuture<String, YourRequestDto, YourResponseDto> future = replyingKafkaTemplate.sendAndReceive(requestRecord); // 阻塞等待响应结果 YourResponseDto response = future.get(5, TimeUnit.SECONDS).value(); return ResponseEntity.ok(response); } }
2. 服务B侧配置与实现
服务B需要监听RequestTopic处理请求,然后把结果发送到ResponseTopic,关键是要保留请求的关联ID(ReplyingKafkaTemplate自动生成的CORRELATION_ID),这样服务A才能正确匹配到对应的响应。
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.ProducerRecord; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Service; import org.apache.kafka.common.header.Header; @Service public class RequestProcessor { private final KafkaTemplate<String, YourResponseDto> kafkaTemplate; public RequestProcessor(KafkaTemplate<String, YourResponseDto> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } @KafkaListener(topics = "RequestTopic") public void processRequest(YourRequestDto request, @Header(KafkaHeaders.CORRELATION_ID) byte[] correlationId) { // 这里写你的业务处理逻辑,生成响应结果 YourResponseDto response = doBusinessLogic(request); // 构造响应消息,设置CORRELATION_ID关联请求 ProducerRecord<String, YourResponseDto> responseRecord = new ProducerRecord<>("ResponseTopic", response); responseRecord.headers().add(KafkaHeaders.CORRELATION_ID, correlationId); // 发送响应到ResponseTopic kafkaTemplate.send(responseRecord); } private YourResponseDto doBusinessLogic(YourRequestDto request) { // 示例业务逻辑,根据实际需求实现 YourResponseDto response = new YourResponseDto(); response.setResult("处理完成"); return response; } }
一些注意事项
- 确保
RequestTopic和ResponseTopic已经提前创建,分区数和副本数根据业务吞吐量配置 - 超时时间要合理设置,既要避免客户端等待过久,也要给服务B足够的处理时间
- 序列化/反序列化要统一,比如服务A和服务B都用JSON序列化,避免消息格式不兼容
- 建议添加异常处理逻辑,比如请求超时、发送失败时返回友好的错误响应给客户端
内容来源于stack exchange
相关产品推荐
相关产品推荐

