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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 10:53:21