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
相关产品推荐
相关产品推荐

