Spring Cloud Stream Kafka自定义错误处理器JSON转换警告修复求助
修复Spring Cloud Stream自定义错误处理器中的JSON转换WARN警告
问题背景
使用Spring Cloud Stream及spring-cloud-stream-binder-kafka 3.2.10版本,配置了自定义错误处理器绑定到consumer1:
spring.cloud.stream.bindings.consumer1-in-0.error-handler-definition=myErrorHandler
自定义错误处理器代码:
@Bean public Consumer<ErrorMessage> myErrorHandler() { return errorMessage -> { System.out.println("Handle error..."); }; }
消费者处理XML格式payload时主动抛出异常:
@Bean public Consumer<String> consumer1() { return msg -> { throw new RuntimeException("Test Error Handling"); }; }
当myErrorHandler执行时,日志中会输出一条WARN信息及堆栈跟踪:
2024-09-17 17:38:17,613 [KafkaConsumerDestination{consumerDestinationName='testTopic', partitions=4, dlqName='null'}.container-0-C-1] WARN org.springframework.cloud.function.context.config.SmartCompositeMessageConverter- Failure during type conversion by org.springframework.cloud.stream.converter.ApplicationJsonMessageMarshallingConverter@d6450a6. Will try the next converter. org.springframework.messaging.converter.MessageConversionException: Could not read JSON: Unrecognized token 'org': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: (String)"org.springframework.messaging.MessageHandlingException: error occurred in message handler [org.springframework.cloud.stream.function.FunctionConfiguration$FunctionToDestinationBinder$1@3d1b0ed0]; nested exception is java.lang.RuntimeException: Test Error Handling, failedMessage=GenericMessage [payload=byte[24595], headers={deliveryAttempt=3, kafka_timestampType=CREATE_TIME, scst_partition=0, kafka_receivedTopic=testTopic, target-protocol=kafka, kafka_offset=29, partitionKey=p0, scst_nativeHeadersPr"[truncated 214 chars]; line: 1, column: 4]; nested exception is com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'org': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: (String)"org.springframework.messaging.MessageHandlingException: error occurred in message handler [org.springframework.cloud.stream.function.FunctionConfiguration$FunctionToDestinationBinder$1@3d1b0ed0]; nested exception is java.lang.RuntimeException: Test Error Handling, failedMessage=GenericMessage [payload=byte[24595], headers={deliveryAttempt=3, kafka_timestampType=CREATE_TIME, scst_partition=0, kafka_receivedTopic=testTopic, target-protocol=kafka, kafka_offset=29, partitionKey=p0, scst_nativeHeadersPr"[truncated 214 chars]; line: 1, column: 4] at org.springframework.messaging.converter.MappingJackson2MessageConverter.convertFromInternal(MappingJackson2MessageConverter.java:237) at org.springframework.cloud.stream.converter.ApplicationJsonMessageMarshallingConverter.convertFromInternal(ApplicationJsonMessageMarshallingConverter.java:115) at org.springframework.messaging.converter.AbstractMessageConverter.fromMessage(AbstractMessageConverter.java:185) at org.springframework.messaging.converter.AbstractMessageConverter.fromMessage(AbstractMessageConverter.java:176) at org.springframework.cloud.function.context.config.SmartCompositeMessageConverter.fromMessage(SmartCompositeMessageConverter.java:63)
问题原因
错误发生时,Spring Cloud Stream会将错误封装为ErrorMessage对象传递给自定义错误处理器。默认的SmartCompositeMessageConverter会按顺序尝试转换器,首先用JSON转换器解析错误信息中的内容,但ErrorMessage包含的异常栈信息是字符串格式而非JSON结构,导致JSON转换失败,触发WARN日志。虽然最终会找到正确的转换器完成处理,但这条不必要的WARN会被输出。
解决方案
方案一:指定错误处理器绑定的Content-Type
为错误处理器的输入绑定指定application/java类型,让转换器直接识别并处理ErrorMessage对象,跳过JSON转换尝试:
spring.cloud.stream.bindings.myErrorHandler-in-0.content-type=application/java
方案二:屏蔽特定WARN日志
如果不需要修改代码或配置逻辑,可通过日志配置将SmartCompositeMessageConverter的WARN级别日志屏蔽。以Logback为例,在logback.xml中添加:
<logger name="org.springframework.cloud.function.context.config.SmartCompositeMessageConverter" level="ERROR"/>
方案三:自定义消息转换器(谨慎使用)
自定义SmartCompositeMessageConverter并移除ApplicationJsonMessageMarshallingConverter,但此方式可能影响其他绑定的消息转换逻辑,需根据实际场景评估后使用。
内容的提问来源于stack exchange,提问作者WoW
相关产品推荐
相关产品推荐

