异步Kafka消息交互问题:App1无法获取App2更新后的对象
问题分析与解决方案
当前的交互逻辑完全不可行,核心问题出在对Kafka异步机制的误解上:
- App1里
kafkaTemplate.send返回的CompletableFuture,仅负责确认消息成功发送到Kafka集群,和App2的处理结果没有任何关联。你在whenCompleteAsync里拿到的result.getProducerRecord().value(),本来就是自己发出去的原始对象,自然看不到App2的更新。 - App2的
@KafkaListener返回CompletableFuture,只是让Listener自身异步执行处理逻辑,并没有提供把更新后的数据回传给App1的通道——你手动创建的SendResult根本没实际发送到Kafka,只是个空壳子。
正确实现思路
要完成「App1发请求→App2处理→返回结果给App1」的同步等待逻辑,必须用请求-响应模式,核心是新增响应主题+用唯一标识关联请求与响应:
步骤1:扩展消息结构与新增主题
给ParsingInfo添加requestId字段(用UUID生成),用来标记每一次请求。同时新增响应主题(比如PARSE_RESPONSE),让App2处理完后把结果发到这个主题。
步骤2:修改App1代码
App1需要完成三件事:生成请求ID、发送请求、监听响应并等待结果
@Autowired KafkaTemplate<String, ParsingInfo> kafkaTemplate; // 线程安全Map,用来关联请求ID和对应的等待Future private final ConcurrentHashMap<String, CompletableFuture<ParsingInfo>> requestFutureMap = new ConcurrentHashMap<>(); // 监听响应主题,收到结果后匹配并完成对应的Future @KafkaListener(topics = Properties.TopicNames.PARSE_RESPONSE, groupId = "app1-response-group") public void handleResponse(ParsingInfo updatedInfo) { CompletableFuture<ParsingInfo> future = requestFutureMap.remove(updatedInfo.getRequestId()); if (future != null) { future.complete(updatedInfo); } } // 发送请求并同步等待响应的方法 public ParsingInfo sendAndWaitForResponse(ParsingInfo parsingInfo) throws Exception { String requestId = UUID.randomUUID().toString(); parsingInfo.setRequestId(requestId); System.out.println("sending message to App2 with requestId: " + requestId); // 确保消息发送成功 kafkaTemplate.send(Properties.TopicNames.PARSE, parsingInfo).join(); // 创建等待Future并放入Map CompletableFuture<ParsingInfo> responseFuture = new CompletableFuture<>(); requestFutureMap.put(requestId, responseFuture); // 等待响应,设置超时时间避免无限阻塞 ParsingInfo newInfo = responseFuture.get(30, TimeUnit.SECONDS); System.out.println("success-get updated data: " + newInfo.getSopSaveData()); return newInfo; }
步骤3:修改App2代码
App2处理完请求后,把更新后的对象发送到响应主题即可:
@Autowired KafkaTemplate<String, ParsingInfo> kafkaTemplate; @KafkaListener(topics = KafkaProperties.TopicNames.PARSE, groupId = KafkaProperties.ConsumerGroupIds.SMART_INTERPRETER) public void listen(ParsingInfo parsingInfo) { System.out.println("Received Message with requestId: " + parsingInfo.getRequestId() + ", data: " + parsingInfo.getSopSaveData()); // 执行更新逻辑 parsingInfo.setSopSaveData("test"); System.out.println("Updated data: " + parsingInfo.getSopSaveData()); // 发送响应到指定主题 kafkaTemplate.send(KafkaProperties.TopicNames.PARSE_RESPONSE, parsingInfo); }
额外注意事项
- 确保
ParsingInfo实现序列化接口(比如Serializable),或配置Kafka使用Json序列化器,否则无法在集群中传输。 - 超时后要及时从
requestFutureMap移除无效条目,避免内存泄漏。 - 可根据业务需求添加异常处理逻辑,比如App2处理失败时返回错误标识,App1捕获后做对应处理。
内容的提问来源于stack exchange,提问作者AstroCoder
相关产品推荐
相关产品推荐

