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

Spring中如何基于Kafka ListenableFuture的onSuccess返回StudentDto给控制器

可行实现方案

由于Kafka消息发送是异步操作,无法在当前方法同步返回StudentDto,必须通过异步响应机制让控制器等待发送结果。以下是几种落地性强的实现方式:

1. 返回CompletableFuture<StudentDto>

将方法返回类型改为CompletableFuture<StudentDto>,通过Future的链式回调处理发送结果,Spring MVC会自动识别并等待Future完成后返回响应:

import org.springframework.util.concurrent.ListenableFuture;
import java.util.concurrent.CompletableFuture;

public class Test {

    public CompletableFuture<StudentDto> publishStudentDto(Student student, String topicName) {
        CompletableFuture<StudentDto> resultFuture = new CompletableFuture<>();
        
        ListenableFuture<SendResult<String, Student>> future = this.studentKafkaTemplate.send(topicName, student);
        future.addCallback(new ListenableFutureCallback<SendResult<String, Student>>() {
            @Override
            public void onSuccess(SendResult<String, Student> result) {
                logger.info("Student created & message published to topic: {} with offset: {} to partition {}", student, result.getRecordMetadata().offset(), result.getRecordMetadata().partition());
                // 转换为StudentDto并完成Future
                StudentDto dto = convertToStudentDto(student);
                resultFuture.complete(dto);
            }

            @Override
            public void onFailure(Throwable ex) {
                logger.error("student not created. Error in publishing student to topic : " + student, ex);
                // 将异常传递给调用方
                resultFuture.completeExceptionally(ex);
            }
        });
        
        return resultFuture;
    }

    // 自行实现Student到StudentDto的转换方法
    private StudentDto convertToStudentDto(Student student) {
        // 转换逻辑
        return new StudentDto();
    }
}

控制器只需直接返回该CompletableFuture<StudentDto>即可,无需额外处理异步逻辑。

2. 使用Spring DeferredResult<StudentDto>

若需要更灵活的超时、异常控制,可借助DeferredResult实现异步响应:

业务层代码

import org.springframework.web.context.request.async.DeferredResult;

public class Test {

    public void publishStudentDto(Student student, String topicName, DeferredResult<StudentDto> deferredResult) {
        ListenableFuture<SendResult<String, Student>> future = this.studentKafkaTemplate.send(topicName, student);
        future.addCallback(new ListenableFutureCallback<SendResult<String, Student>>() {
            @Override
            public void onSuccess(SendResult<String, Student> result) {
                logger.info("Student created & message published to topic: {} with offset: {} to partition {}", student, result.getRecordMetadata().offset(), result.getRecordMetadata().partition());
                StudentDto dto = convertToStudentDto(student);
                deferredResult.setResult(dto);
            }

            @Override
            public void onFailure(Throwable ex) {
                logger.error("student not created. Error in publishing student to topic : " + student, ex);
                deferredResult.setErrorResult(ex);
            }
        });
    }

    private StudentDto convertToStudentDto(Student student) {
        return new StudentDto();
    }
}

控制器代码

import org.springframework.web.context.request.async.DeferredResult;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class StudentController {

    private final Test testService;

    public StudentController(Test testService) {
        this.testService = testService;
    }

    @PostMapping("/publish-student")
    public DeferredResult<StudentDto> publishStudent(@RequestBody Student student) {
        DeferredResult<StudentDto> deferredResult = new DeferredResult<>(5000L); // 设置5秒超时
        deferredResult.onTimeout(() -> deferredResult.setErrorResult("消息发送超时"));
        
        testService.publishStudentDto(student, "student-topic", deferredResult);
        return deferredResult;
    }
}

3. 同步等待发送结果(不推荐)

如果业务场景允许阻塞线程,可直接调用ListenableFuture的get()方法同步等待结果,但会降低系统并发能力,生产环境慎用:

public StudentDto publishStudentDto(Student student, String topicName) throws ExecutionException, InterruptedException {
    ListenableFuture<SendResult<String, Student>> future = this.studentKafkaTemplate.send(topicName, student);
    // 同步阻塞等待发送结果
    SendResult<String, Student> result = future.get();
    
    logger.info("Student created & message published to topic: {} with offset: {} to partition {}", student, result.getRecordMetadata().offset(), result.getRecordMetadata().partition());
    return convertToStudentDto(student);
}

需注意捕获ExecutionException(封装发送异常)和InterruptedException(线程中断异常)。


内容的提问来源于stack exchange,提问作者PTInnovating

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 11:10:20