You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.08 12:53:01