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
推测是消息头不匹配导致问题,咨询以下两点:
- 使用Spring Kafka回复模板时,可接受的回复Topic头是什么?
- 双方通信的消息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
相关产品推荐
相关产品推荐

