如何实现Kafka消费者结果回传?微服务场景下的解决方案
问题描述
我正在开发一套以Kafka作为消息中间件的微服务应用。在user-service中,存在一个控制器会调用函数向payment-service发送请求。目前我的Kafka生产者与消费者配置均正常运行,但我希望了解是否能够将payment-service的响应结果返回至user-service的控制器,若可行,该如何实现此功能?
现有代码
user-service 业务方法
public List<Donation> getUserDonationHistory(String userId) { // 1. Make a request to the payment-service from which to get the donation-history producerService.requestUserDonationHistory(userId); return null; }
user-service Kafka生产者
public void requestUserDonationHistory(String userId) { log.info("Requesting the donation-history for the user with the ID: {}", userId); kafkaTemplate.send("request-donation-history", userId); }
payment-service Kafka消费者
@KafkaListener(topics = "request-donation-history", groupId = "groupId") public void getDonationHistoryRequest(String userId) { log.info("Got a donationHistoryRequest. Raw userId: {}", userId); Long parsedUserId = Long.parseLong(gson.fromJson(userId, String.class)); log.info("Parsed userId: {}", parsedUserId); producerService.sendUserDonationHistory(parsedUserId); }
可行性与实现方案
完全可以实现,核心是通过Kafka请求-响应模式,用唯一请求ID关联请求与响应,配合异步Future实现控制器的结果返回。以下是具体实现步骤:
1. 定义带唯一标识的请求/响应DTO
为了关联请求和响应,需要给每个请求添加唯一ID,不再直接传递userId:
// 请求DTO public class DonationHistoryRequest { private String requestId; private String userId; // 构造器、getter、setter } // 响应DTO public class DonationHistoryResponse { private String requestId; private List<Donation> donations; // 构造器、getter、setter }
2. 改造user-service的生产者逻辑
生成唯一请求ID,用CompletableFuture存储请求上下文,等待响应完成:
@Service public class KafkaProducerService { private final KafkaTemplate<String, Object> kafkaTemplate; // 存储请求ID与对应Future的映射 private final ConcurrentHashMap<String, CompletableFuture<List<Donation>>> futureMap = new ConcurrentHashMap<>(); public KafkaProducerService(KafkaTemplate<String, Object> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public CompletableFuture<List<Donation>> requestUserDonationHistory(String userId) { String requestId = UUID.randomUUID().toString(); DonationHistoryRequest request = new DonationHistoryRequest(requestId, userId); log.info("请求用户{}的捐赠历史,请求ID:{}", userId, requestId); kafkaTemplate.send("request-donation-history", request); CompletableFuture<List<Donation>> future = new CompletableFuture<>(); futureMap.put(requestId, future); return future; } // 响应回调:完成Future public void completeRequest(String requestId, List<Donation> donations) { CompletableFuture<List<Donation>> future = futureMap.remove(requestId); if (future != null) { future.complete(donations); } } // 异常回调:标记请求失败 public void failRequest(String requestId, Throwable ex) { CompletableFuture<List<Donation>> future = futureMap.remove(requestId); if (future != null) { future.completeExceptionally(ex); } } }
同时修改业务方法,等待响应结果(设置超时避免无限阻塞):
public List<Donation> getUserDonationHistory(String userId) throws Exception { // 等待响应,超时时间设为5秒 return producerService.requestUserDonationHistory(userId) .get(5, TimeUnit.SECONDS); }
3. 在user-service添加响应消费者
监听payment-service返回的响应主题,触发Future完成:
@KafkaListener(topics = "response-donation-history", groupId = "user-service-group") public void handleDonationHistoryResponse(DonationHistoryResponse response) { log.info("收到请求ID{}的捐赠历史响应", response.getRequestId()); producerService.completeRequest(response.getRequestId(), response.getDonations()); }
4. 改造payment-service的处理逻辑
接收带请求ID的消息,处理后返回带相同ID的响应:
@KafkaListener(topics = "request-donation-history", groupId = "payment-service-group") public void getDonationHistoryRequest(DonationHistoryRequest request) { log.info("收到捐赠历史请求,用户ID:{},请求ID:{}", request.getUserId(), request.getRequestId()); Long parsedUserId = Long.parseLong(request.getUserId()); // 调用业务逻辑获取捐赠历史 List<Donation> donations = donationService.getDonationsByUserId(parsedUserId); // 构造响应并发送到响应主题 DonationHistoryResponse response = new DonationHistoryResponse(request.getRequestId(), donations); producerService.sendDonationHistoryResponse(response); }
payment-service的响应生产者方法:
public void sendDonationHistoryResponse(DonationHistoryResponse response) { log.info("发送请求ID{}的捐赠历史响应", response.getRequestId()); kafkaTemplate.send("response-donation-history", response); }
5. 关键注意事项
- 超时处理:必须设置Future的等待超时,避免线程长期阻塞。
- 序列化配置:确保Kafka使用Jackson等序列化器处理DTO,避免消息格式错误。
- 异常处理:payment-service处理失败时,要发送错误响应或触发user-service的
failRequest方法。 - 幂等性:payment-service的业务逻辑要保证幂等,避免重复消费消息导致数据重复。
内容的提问来源于stack exchange,提问作者Marius Carchilan
相关产品推荐
相关产品推荐

