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

Spring Kafka回复模板与Kafkajs交互报错及适配问题咨询

问题描述

使用Kafkajs向Spring Boot实现的Kafka回复模板应用发送消息时,Spring端抛出异常,相关代码及错误信息如下:

Spring端消费代码

@KafkaListener(topics = KHQRTopic.GENERATE_KHQR, groupId = KafkaConfiguration.CONSUMER_GROUP)
@SendTo
public Message<?> onEventGenerateKHQR(ConsumerRecord<String, KHQRRequest> consumerRecord) {
        log.info("Received event {}", consumerRecord.value());
...
}

Node.js(Kafkajs)发送代码

export const getResult = async <T>(messageOption: MessageOption, data: any) => {
  const processID = new Date().getTime().toString()
  let payload: string | Buffer = ''
  
    if (messageOption.avroSchemaName) {
      payload = await registry.encode(
        getSchemaId(messageOption.avroSchemaName),
        data
      );
    } else {
      payload = JSON.stringify(data)
    }
    
  
    const replyTopic = messageOption.replyTopic ?? `${messageOption.sendTopic}.reply`
    console.log('wait result form', replyTopic)
    await producer.send({
      topic: messageOption.sendTopic,
      messages: [
        { key: processID, value: payload,
          headers: {
          'reply-topic': replyTopic
      } }],
      
    });

    return await getMessageResponse<T>(
      replyTopic,
      processID,
      messageOption.timer
    );
}

Spring端抛出的错误

org.springframework.kafka.listener.ListenerExecutionFailedException: Listener failed
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.decorateException(KafkaMessageListenerContainer.java:2946) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2887) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeOnMessage(KafkaMessageListenerContainer.java:2854) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.lambda$doInvokeRecordListener$57(KafkaMessageListenerContainer.java:2772) ~[spring-kafka-3.0.12.jar:3.0.12]
    at io.micrometer.observation.Observation.lambda$observe$4(Observation.java:544) ~[micrometer-observation-1.11.5.jar:1.11.5]
    at io.micrometer.observation.Observation.observeWithContext(Observation.java:603) ~[micrometer-observation-1.11.5.jar:1.11.5]
    at io.micrometer.observation.Observation.observe(Observation.java:544) ~[micrometer-observation-1.11.5.jar:1.11.5]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:2770) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:2622) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2508) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:2150) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1505) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1469) ~[spring-kafka-3.0.12.jar:3.0.12]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1344) ~[spring-kafka-3.0.12.jar:3.0.12]
    at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]
Caused by: java.lang.IllegalStateException: With no topic header, a defaultTopic is required

推测是消息头不匹配导致问题,咨询以下两点:

  1. 使用Spring Kafka回复模板时,可接受的回复Topic头是什么?
  2. 双方通信的消息Payload需满足什么格式要求?

解答

1. Spring Kafka回复模板接受的回复Topic头

Spring Kafka的@SendTo注解在处理请求/回复模式时,默认读取**kafka_replyTopic**这个消息头(注意键名的精确大小写),而非你代码中的reply-topic。

错误日志里的With no topic header, a defaultTopic is required,就是因为Spring找不到指定的回复头,且未配置默认回复Topic导致的。

修改Node.js代码的headers部分,将键名改为kafka_replyTopic即可:

headers: {
  'kafka_replyTopic': replyTopic
}

如果需要自定义头名称,也可以通过配置KafkaListenerContainerFactory实现,但推荐使用默认值减少配置复杂度。

2. 消息Payload格式要求

双方的Payload格式必须保持一致,具体取决于Spring端的反序列化配置和Node.js端的序列化方式:

情况一:JSON格式

  • Node.js端:直接用JSON.stringify(data)生成标准JSON字符串即可。
  • Spring端:需配置对应的JSON反序列化器——要么用StringDeserializer接收值后手动反序列化为KHQRRequest;要么配置JsonDeserializer作为值反序列化器,并指定目标类型为KHQRRequest。

情况二:Avro格式

  • Node.js端:通过Schema Registry将数据编码为Avro二进制,确保使用的Schema与Spring端完全一致。
  • Spring端:需配置KafkaAvroDeserializer作为值反序列化器,指定Schema Registry地址,同时保证KHQRRequest类的结构与Avro Schema完全匹配(可通过Avro工具自动生成Java类)。

另外,消息key的序列化/反序列化也要保持一致,比如双方都使用StringSerializer/StringDeserializer。


内容的提问来源于stack exchange,提问作者phen pani

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 04:35:54