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

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,提问作者임띵동

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 19:58:08