Kafka Stream消费端遭遇JsonParseException转换异常求助
问题场景
配置了如下Spring Cloud Stream Kafka消费端配置:
spring: cloud: stream: function: definition: handleCatalogEvent bindings: handleCatalogEvent-in-0: content-type: application/json destination: catalog_change group: back-group consumer: configuration: json.value.type: com.test.domain.ChangeNotification json.fail.invalid.schema: true useNativeEncoding: true kafka: binder: brokers: kafka-dev.test.priv:7887 auto-create-topics: false consumer-properties: auto.offset.reset: latest auto.commit.interval.ms: 1000 specific.avro.reader: true auto.register.schemas: false key.deserializer: org.apache.kafka.common.serialization.StringDeserializer value.deserializer: io.confluent.kafka.serializers.json.KafkaJsonSchemaDeserializer schema.registry.url: https://schema-registry.dev.test.priv:8787
对应的消费处理Bean:
@Bean public Consumer<Message<ChangeNotification>> handleCatalogEvent() { return event -> { log.info("- - - - - - - - - - - - - - - - - - - Nature Of Change : {} - - - - - - - - - - - - - - - - -", event.getPayload().getNatureOfChange()); log.info("- - - - - - - - - - - - - - - - - - - - - - PAYLOAD : {} - - - - - - - - - - - - - - - - -", event.getPayload().getItem()); }; }
消费时抛出JSON解析异常:
Caused by: com.fasterxml.jackson.core.JsonParseException: Unexpected character ('i' (code 105)): was expecting double-quote to start field name at [Source: (String)"{idCatalog=5c824fc5b0efe060e87f056b, code=008, codeContext=[01], label=Spécialité PS, natureOfChange=INSERT, item={id=6373b315d9428a128d0ffc6c, code=kn17, label=string, startEffectiveDate=2022-11-15T10:25:23.21Z}}"; line: 1, column: 3] at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2391) ~[jackson-core-2.13.1.jar:2.13.1] at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:735) ~[jackson-core-2.13.1.jar:2.13.1] at com.fasterxml.jackson.core.base.ParserMinimalBase._reportUnexpectedChar(ParserMinimalBase.java:659) ~[jackson-core-2.13.1.jar:2.13.1] at com.fasterxml.jackson.core.json.ReaderBasedJsonParser._handleOddName(ReaderBasedJsonParser.java:1860) ~[jackson-core-2.13.1.jar:2.13.1] at com.fasterxml.jackson.core.json.ReaderBasedJsonParser.nextToken(ReaderBasedJsonParser.java:734) ~[jackson-core-2.13.1.jar:2.13.1] at com.fasterxml.jackson.databind.deser.BeanDeserializer.deserialize(BeanDeserializer.java:176) ~[jackson-databind-2.13.1.jar:2.13.1] at com.fasterxml.jackson.databind.deser.DefaultDeserializationContext.readRootValue(DefaultDeserializationContext.java:322) ~[jackson-databind-2.13.1.jar:2.13.1] at com.fasterxml.jackson.databind.ObjectMapper._readMapAndClose(ObjectMapper.java:4674) ~[jackson-databind-2.13.1.jar:2.13.1] at com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:3629) ~[jackson-databind-2.13.1.jar:2.13.1] at org.springframework.messaging.converter.MappingJackson2MessageConverter.convertFromInternal(MappingJackson2MessageConverter.java:232) ~[spring-messaging-5.3.23.jar:5.3.23] ... 45 common frames omitted
异常原因
- 数据格式不匹配:从错误日志的源字符串可见,消息内容是
{idCatalog=5c824fc5b0efe060e87f056b,...},这是非标准JSON格式——字段名无引号、键值用=分隔,属于Java对象toString()输出或Groovy Map格式,无法被Jackson解析为标准JSON。 - 配置冲突:当前配置同时启用了两种反序列化逻辑:
content-type: application/json触发Spring Cloud Stream的MappingJackson2MessageConverter尝试解析消息;- 虽指定了Confluent的
KafkaJsonSchemaDeserializer,但useNativeEncoding: true未正确生效,导致Spring的转换器优先处理了非标准格式的消息。
解决方法
方案一:修正生产者发送标准JSON数据(推荐)
让生产者发送符合RFC标准的JSON格式数据(字段名带双引号、键值用:分隔),示例:
{"idCatalog":"5c824fc5b0efe060e87f056b", "code":"008", "codeContext":["01"], "label":"Spécialité PS", "natureOfChange":"INSERT", "item":{"id":"6373b315d9428a128d0ffc6c", "code":"kn17", "label":"string", "startEffectiveDate":"2022-11-15T10:25:23.21Z"}}
此情况下,现有消费端的Jackson转换器可正常解析数据。
方案二:调整消费端配置适配非标准格式
若无法修改生产者,可通过以下方式调整:
禁用Spring消息转换器,启用原生Kafka反序列化:
移除content-type: application/json配置,确保useNativeEncoding: true生效,让Confluent的KafkaJsonSchemaDeserializer直接处理消息。注意:该反序列化器仅支持符合JSON Schema规范的标准JSON,若生产者发送非标准格式则不适用。自定义消息转换器:
编写自定义MessageConverter解析非标准格式字符串,转换为ChangeNotification对象:@Bean public MessageConverter customMessageConverter() { return new AbstractMessageConverter(MimeType.valueOf("text/plain")) { @Override protected Object convertFromInternal(Message<?> message, Class<?> targetClass, Object conversionHint) { String payload = new String((byte[]) message.getPayload(), StandardCharsets.UTF_8); // 将非标准格式转换为标准JSON String normalizedJson = payload.replaceAll("(\\w+)=([^,}]+)", "\"$1\":\"$2\"") .replaceAll("\\[([^\\]]+)\\]", "\"[$1]\""); try { return new ObjectMapper().readValue(normalizedJson, targetClass); } catch (JsonProcessingException e) { throw new RuntimeException("解析自定义格式失败", e); } } }; }同时修改绑定配置的
content-type为text/plain,让自定义转换器生效。
方案三:修正Confluent JSON Schema反序列化配置(若生产者用Schema Registry)
若生产者通过Confluent Schema Registry发送JSON Schema格式消息,需修正消费端配置:
- 移除绑定配置中
consumer.configuration下的json.value.type,改用spring.cloud.stream.kafka.binder.consumer-properties.json.value.type; - 将
specific.avro.reader: true改为specific.json.reader: true(适配JSON Schema); - 移除
content-type: application/json,保留useNativeEncoding: true。
修正后的核心配置片段:
spring: cloud: stream: bindings: handleCatalogEvent-in-0: destination: catalog_change group: back-group consumer: useNativeEncoding: true kafka: binder: consumer-properties: specific.json.reader: true json.value.type: com.test.domain.ChangeNotification # 其他原有配置...
内容的提问来源于stack exchange,提问作者Aymen Kanzari

