【问题标题】:Kafka Consumer Error卡夫卡消费者错误
【发布时间】:2018-12-04 00:25:41
【问题描述】:

我正在使用一个 kafka 生产者和一个 Spring kafka 消费者。我正在使用 Json 序列化器和反序列化器。每当我尝试从主题中读取消费者中的消息时,我都会收到以下错误:

org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition fan_topic-0 at offset 154. If needed, please seek past the record to continue consumption.
Caused by: java.lang.IllegalStateException: No type information in headers and no default type provided

我没有在生产者和消费者中配置任何有关标头的内容。我在这里错过了什么?

【问题讨论】:

    标签: apache-kafka kafka-consumer-api kafka-producer-api spring-kafka


    【解决方案1】:

    我相信您错过了这样一个事实,即必须在 ConsumerFactory 上配置 JsonDeserializer 并使用适当的默认类型进行反序列化,而不是在 Kafka 属性中。

    所有信息都显示在文档中:https://docs.spring.io/spring-kafka/docs/2.1.7.RELEASE/reference/html/_reference.html#serdes

    【讨论】:

    • 在属性中添加 JsonDeserializer.VALUE_DEFAULT_TYPE 解决了这个问题。
    • 如果我使用 EmbeddedKafka 做一个简单的测试,所有配置都放在 yaml 文件中?在这种情况下,CosumerFactory bean 全部由 yaml 构建,而不是由我构建;我不想在 config 中重复 yaml 中的所有信息并将其传递给 consumerFactory 的构造函数以提供类型信息。
    【解决方案2】:

    只是添加到上面的答案,

    以下更改为我解决了。

    config.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);
    

    添加

    return new DefaultKafkaConsumerFactory<>(config, new StringDeserializer(), new JsonDeserializer<>(String.class));
    

    而不是

    return new DefaultKafkaConsumerFactory<String, String>(config);
    

    供参考,

    deserialize 中的以下方法期待标题和“Assert.state..”抛出IllegalStateException

     @Override
            public T deserialize(String topic, Headers headers, byte[] data) {
                JavaType javaType = this.typeMapper.toJavaType(headers);
                if (javaType == null) {
                    Assert.state(this.targetType != null, "No type information in headers and no default type provided");
                    return deserialize(topic, data);
                }
                else {
                    try {
                        return this.objectMapper.readerFor(javaType).readValue(data);
                    }
                    catch (IOException e) {
                        throw new SerializationException("Can't deserialize data [" + Arrays.toString(data) +
                                "] from topic [" + topic + "]", e);
                    }
                }
            }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-07-03
      • 2018-05-05
      • 2021-08-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-10-28
      • 2015-12-18
      相关资源
      最近更新 更多