Java服务通过ReplyingKafkaTemplate与Kotlin服务通信收不到消息
核心问题点
1. ReplyingKafkaTemplate实例配置丢失
Java端配置中创建ReplyingKafkaTemplate后,已经设置了超时时间,但return时重新new了一个全新实例,导致之前的超时配置完全失效,同时可能引发请求响应流程中的异常。
2. 序列化与反序列化不匹配
Java端生产者使用JsonSerializer直接序列化KafkaResponseDto对象,但Kotlin端消费者用StringDeserializer接收字符串后手动反序列化,存在两个潜在问题:
JsonSerializer默认会在消息头中添加__TypeId__类型标识,StringDeserializer不会处理该标识,可能导致反序列化后的字符串包含额外无效信息,引发ObjectMapper解析失败;- Java端发送的
KafkaResponseDto与Kotlin端要解析的KafkaMemberValidateRequestDto如果字段结构、名称(含大小写)不匹配,会直接触发解析异常,最终返回null。
3. Kotlin端ObjectMapper配置缺失
如果未给Kotlin端的ObjectMapper注册Kotlin模块,Jackson无法正确解析Kotlin数据类,会抛出反序列化异常,导致方法返回null。
修复步骤
1. 修复ReplyingKafkaTemplate创建代码
返回已配置超时的实例,而非重新创建:
@Bean public ReplyingKafkaTemplate<String, Object, String> replyKafkaTemplate (ProducerFactory<String, Object> pf, KafkaMessageListenerContainer<String, String> container) { ReplyingKafkaTemplate<String, Object, String> replyingKafkaTemplate = new ReplyingKafkaTemplate<>(pf, container); replyingKafkaTemplate.setDefaultReplyTimeout(Duration.ofMillis(5000)); return replyingKafkaTemplate; // 返回已配置的实例,而非重新new }
2. 统一序列化逻辑(二选一)
方案A:Java端手动序列化JSON字符串发送
将KafkaResponseDto提前序列化为JSON字符串,改用StringSerializer发送,确保Kotlin端接收的是标准JSON:
// Java端replyRecord方法修改 public Object replyRecord(KafkaResponseDto requestData) throws ExecutionException, InterruptedException, JsonProcessingException, TimeoutException { ObjectMapper objectMapper = new ObjectMapper(); String jsonRequest = objectMapper.writeValueAsString(requestData); ProducerRecord<String, Object> record = new ProducerRecord<>(kafkaRestApiTopic, jsonRequest); record.headers().add(new RecordHeader(KafkaHeaders.REPLY_TOPIC, kafkaRestApiTopic.getBytes())); RequestReplyFuture<String, Object, String> sendAndReceive = replyingKafkaTemplate.sendAndReceive(record); SendResult<String, Object> sendResult = sendAndReceive.getSendFuture().get(10, TimeUnit.SECONDS); sendResult.getProducerRecord().headers().forEach(header -> System.out.println(header.key() + ":" + new String(header.value()))); ConsumerRecord<String, String> consumerRecord = sendAndReceive.get(); return consumerRecord.value(); }
同时修改Java端生产者配置的序列化器:
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
方案B:Kotlin端直接用JsonDeserializer接收DTO
让Kotlin端消费者直接反序列化为目标DTO,避免手动解析:
// Kotlin端监听方法修改 @KafkaListener(topics = ["\${devsoo.kafka.restapi.topic}"], groupId = "member-server-consumer-group") @SendTo fun getMemberRequestDto(message: KafkaMemberValidateRequestDto): String { val member = memberService.getInstance(message) return objectMapper.writeValueAsString(member.toKafkaResponseDto()) }
同时在Kotlin端的Kafka配置中设置:
// 消费者配置示例 props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = JsonDeserializer::class.java.name props[JsonDeserializer.TRUSTED_PACKAGES] = "*" props[JsonDeserializer.VALUE_DEFAULT_TYPE] = KafkaMemberValidateRequestDto::class.java.name
3. 完善Kotlin端ObjectMapper配置
添加Jackson Kotlin模块,确保能正确解析Kotlin数据类:
val objectMapper = ObjectMapper().registerModule(KotlinModule())
同时确保KafkaResponseDto(Java)与KafkaMemberValidateRequestDto(Kotlin)的字段名称、类型完全匹配(注意Jackson默认区分大小写,可配置ACCEPT_CASE_INSENSITIVE_PROPERTIES兼容)。
4. 临时移除ErrorHandlingDeserializer定位问题
如果问题仍未解决,暂时移除ErrorHandlingDeserializer,直接使用底层反序列化器,这样反序列化失败时会抛出明确异常,便于快速定位问题根源。
内容的提问来源于stack exchange,提问作者임띵동

