【问题标题】:Apache Kafka: ...StringDeserializer is not an instance of ...DeserializerApache Kafka:...StringDeserializer 不是 ...Deserializer 的实例
【发布时间】:2018-07-29 13:29:31
【问题描述】:

在我的简单应用程序中,我试图实例化一个 KafkaConsumer,我的代码几乎是 code from javadoc 的副本(“自动偏移提交”):

@Slf4j
public class MyKafkaConsumer {

    public MyKafkaConsumer() {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "test");
        props.put("enable.auto.commit", "true");
        props.put("auto.commit.interval.ms", "1000");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
        consumer.subscribe( Arrays.asList("mytopic"));
        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(100);
            for (ConsumerRecord<String, String> record : records)
                log.info( record.offset() + record.key() + record.value() );
                //System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
        }
    }
}

如果我尝试实例化它,我会得到:

org.apache.kafka.common.KafkaException: Failed to construct kafka consumer
        at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:781)
        at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:635)
        at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:617)
at ...MyKafkaConsumer.<init>(SikomKafkaConsumer.java:23)
    ...
    Caused by: org.apache.kafka.common.KafkaException: org.apache.kafka.common.serialization.StringDeserializer is not an instance of org.apache.kafka.common.serialization.Deserializer
        at org.apache.kafka.common.config.AbstractConfig.getConfiguredInstance(AbstractConfig.java:248)
        at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:680)
        ... 48 more

如何解决这个问题?

【问题讨论】:

  • 你是怎么解决这个问题的?
  • 我不确定,但现在工作的代码与我的问题相同,但可能是依赖问题:我有这个:'org.apache.kafka:kafka-客户:1.0.0'
  • 谢谢。对我来说也是传递依赖的问题
  • 你可以添加这个作为答案,我会接受

标签: apache-kafka kafka-consumer-api


【解决方案1】:

你的自定义类需要实现,org.apache.kafka.common.serialization.Deserializer。

喜欢

import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.common.serialization.Deserializer;
import org.codehaus.jackson.map.ObjectMapper;

import java.io.Serializable;
import java.util.Map;

//Developed by Arun Singh
public class Employee implements Serializable, Serializer, **Deserializer** {

@Override
    public Object deserialize(String s, byte[] bytes) {
        ObjectMapper mapper = new ObjectMapper();
        Employee employee = null;
        try {
            //employee = mapper.readValue(bytes, Employee.class);
            employee = mapper.readValue(bytes.toString(), Employee.class);
        } catch (Exception e) {
            e.printStackTrace();
        }
        return employee;
    }

    @Override
    public Object deserialize(String topic, Headers headers, byte[] data) {
        ObjectMapper mapper = new ObjectMapper();
        Employee employee = null;
        try {
            //employee = mapper.readValue(bytes, Employee.class);
            employee = mapper.readValue(data.toString(), Employee.class);
        } catch (Exception e) {
            e.printStackTrace();
        }
        return employee;
    }

    public void close() {

    }
}

【讨论】:

    【解决方案2】:

    这可能是 Kafka 类加载的问题。
    将类加载器设置为 null 可能会有所帮助。

    ...
    Thread currentThread = Thread.currentThread();    
    ClassLoader savedClassLoader = currentThread.getContextClassLoader();
    
    currentThread.setContextClassLoader(null);
    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
    
    currentThread.setContextClassLoader(savedClassLoader);
    ...
    

    有完整解释:
    https://stackoverflow.com/a/50981469/1673775

    【讨论】:

      【解决方案3】:

      不确定这是否最终解决了您的错误,但请注意,在将 spring-kafka-test(版本 2.1.x,从版本 2.1.5 开始)与 1.1.x kafka-clients jar 一起使用时,您需要覆盖某些传递依赖项,如下所示:

      <dependency>
          <groupId>org.springframework.kafka</groupId>
          <artifactId>spring-kafka</artifactId>
          <version>${spring.kafka.version}</version>
      </dependency>
      
      <dependency>
          <groupId>org.springframework.kafka</groupId>
          <artifactId>spring-kafka-test</artifactId>
          <version>${spring.kafka.version}</version>
          <scope>test</scope>
      </dependency>
      
      <dependency>
          <groupId>org.apache.kafka</groupId>
          <artifactId>kafka-clients</artifactId>
          <version>1.1.1</version>
      </dependency>
      
      <dependency>
          <groupId>org.apache.kafka</groupId>
          <artifactId>kafka-clients</artifactId>
          <version>1.1.1</version>
          <classifier>test</classifier>
      </dependency>
      
      <dependency>
          <groupId>org.apache.kafka</groupId>
          <artifactId>kafka_2.11</artifactId>
          <version>1.1.1</version>
          <scope>test</scope>
      </dependency>
      
      <dependency>
          <groupId>org.apache.kafka</groupId>
          <artifactId>kafka_2.11</artifactId>
          <version>1.1.1</version>
          <classifier>test</classifier>
          <scope>test</scope>
      </dependency>
      

      所以你的传递依赖肯定是个问题

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2021-08-03
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-10-11
        相关资源
        最近更新 更多