【问题标题】:Could not decode json type in spring cloud stream DefaultKafkaHeaderMapper无法解码 Spring Cloud Stream DefaultKafkaHeaderMapper 中的 json 类型
【发布时间】:2020-12-14 07:11:01
【问题描述】:

我们正在使用 spring-cloud-stream 并计划升级我们的 Kafka 版本。
我们的应用程序使用 spring-cloud-stream:2.0.0 (spring-kafka 2.1.7) 和 apache kafka 服务器 1.0.1 并且还使用spring-cloud-sleuth:2.0.0 进行跟踪。
我们将把我们的 Kafka 服务器升级到版本 2.3.0,所以它需要升级到 spring-boot 2.2.x (Hoxton)spring-cloud-sleuth:2.2.0spring-cloud-stream:3.0.3 (Horsham.SR3)
我们有大约 200 个使用 Kafka 的应用程序,因此升级将逐渐进行,因此作为中间状态,我们将在较新的版本上拥有 生产者,而在旧版本上拥有 消费者。 我们的消费者正在使用@StreamListener

在我们的测试中,我们遇到了一个问题:解析大多数类型为 String 的标头并获得以下信息:

ERROR 27448 --- [container-0-C-1] o.s.c.s.b.k.KafkaMessageChannelBinder$4  : Could not decode json type: ecb89ccb3e79418b for key: X-B3-TraceId
com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'ecb89ccb3e79418b': was expecting ('true', 'false' or 'null')
 at [Source: (byte[])"ecb89ccb3e79418b"; line: 1, column: 33]
    at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:1804) ~[jackson-core-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:679) ~[jackson-core-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3526) ~[jackson-core-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._handleUnexpectedValue(UTF8StreamJsonParser.java:2621) ~[jackson-core-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:826) ~[jackson-core-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:723) ~[jackson-core-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.databind.ObjectMapper._initForReading(ObjectMapper.java:4141) ~[jackson-databind-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.databind.ObjectMapper._readMapAndClose(ObjectMapper.java:4000) ~[jackson-databind-2.9.6.jar:2.9.6]
    at com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:3091) ~[jackson-databind-2.9.6.jar:2.9.6]
    at org.springframework.kafka.support.DefaultKafkaHeaderMapper.lambda$toHeaders$1(DefaultKafkaHeaderMapper.java:233) ~[spring-kafka-2.1.7.RELEASE.jar:2.1.7.RELEASE]
    at java.lang.Iterable.forEach(Iterable.java:75) ~[na:1.8.0_221]
    at org.springframework.kafka.support.DefaultKafkaHeaderMapper.toHeaders(DefaultKafkaHeaderMapper.java:216) ~[spring-kafka-2.1.7.RELEASE.jar:2.1.7.RELEASE]
    at org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder$4.toHeaders(KafkaMessageChannelBinder.java:554) ~[spring-cloud-stream-binder-kafka-2.0.0.RELEASE.jar:2.0.0.RELEASE]
    at org.springframework.kafka.support.converter.MessagingMessageConverter.toMessage(MessagingMessageConverter.java:106) ~[spring-kafka-2.1.7.RELEASE.jar:2.1.7.RELEASE]
    at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.toMessagingMessage(MessagingMessageListenerAdapter.java:229) ~[spring-kafka-2.1.7.RELEASE.jar:2.1.7.RELEASE]
...

虽然类型标题是:

{spanTraceId=java.lang.String, spanId=java.lang.String, spanParentSpanId=java.lang.String, nativeHeaders=org.springframework.util.LinkedMultiValueMap, X-B3-SpanId=java.lang.String, X-B3-ParentSpanId=java.lang.String, scst_partition=java.lang.Integer, X-B3-Sampled=java.lang.String, X-B3-TraceId=java.lang.String, spanSampled=java.lang.String, contentType=java.lang.String}

例如,Sleuth 添加的X-B3-SpanId 是String 类型,其值为:ecb89ccb3e79418b,它不是JSON 字符串,因此ObjectMapper fails 转换到这里的字符串对象:

headers.put(h.key(), getObjectMapper().readValue(h.value(), type))

看起来当我们有 String 类型时它不应该使用 ObjectMapper,因此我们的旧消费者失败了。

在使用新的生产者和旧的消费者时,有没有办法防止这个问题?

【问题讨论】:

    标签: apache-kafka spring-kafka spring-cloud-stream spring-cloud-sleuth spring-cloud-stream-binder-kafka


    【解决方案1】:

    您可以将DefaultKafkaHeaderMapper 配置为与旧版本兼容:

        /**
         * Set to true to encode String-valued headers as JSON ("..."), by default just the
         * raw String value is converted to a byte array using the configured charset. Set to
         * true if a consumer of the outbound record is using Spring for Apache Kafka version
         * less than 2.3
         * @param encodeStrings true to encode (default false).
         * @since 2.3
         */
        public void setEncodeStrings(boolean encodeStrings) {
            this.encodeStrings = encodeStrings;
        }
    
    

    另见https://docs.spring.io/spring-cloud-stream-binder-kafka/docs/3.0.10.RELEASE/reference/html/spring-cloud-stream-binder-kafka.html#_kafka_binder_properties

    spring.cloud.stream.kafka.binder.headerMapperBeanName

    【讨论】:

    • 嗨@gary,这个建议不起作用。我们尝试使用 Hoxton.SR3 进行生产,并按照建议在 DefaultKafkaHeaderMapper 上设置了 encodeStrings flag = true。使用此设置,消费者能够处理所有 String 标头,但在内容类型标头上失败。此标头的类型值已更改为 org.springframework.util.MimeType(由于标志更改),值为“application/json”,因此它试图将其转换为 JSON 并在此处失败:headers.put(h.key(), getObjectMapper().readValue(h.value(), type)); with NPE。
    • 尝试使用带有属性集的BinderHeaderMapper,而不是DefaultKafkaHeaderMapper;它最初是一个克隆,但我在那里看到了一些额外的代码来处理MimeType
    • 谢谢@gary - BinderHeaderMapper 做到了,从spring-kafka:2.3.7开始没有其他与标​​题相关的问题
    猜你喜欢
    • 2020-03-28
    • 2020-08-22
    • 1970-01-01
    • 2019-02-07
    • 1970-01-01
    • 2015-08-27
    • 2020-09-27
    • 2023-04-05
    • 2020-03-24
    相关资源
    最近更新 更多