【问题标题】:Json Message Communication With Apache Kafka in Spring BootSpring Boot 中与 Apache Kafka 的 Json 消息通信
【发布时间】:2020-10-14 12:21:38
【问题描述】:

我已经实现了两个不同的应用程序,一个用于生产者,另一个用于消费者,并且消息已在 apache kafka 的支持下传递。当我发布字符串消息时,通信正常完成,但是当我传递 Json 消息时,发生以下错误。

错误

java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer

Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition Kafka_Example_Json-0 at offset 0. If needed, please seek past the record to continue consumption.

Caused by: java.lang.IllegalArgumentException: The class 'com.benz.kafka.api.model.User' is not in the trusted packages: [java.util, java.lang, com.benz.kafka.consumer.api.model, com.benz.kafka.consumer.api.model.*]. If you believe this class is safe to deserialize, please provide its name. If the serialization is only done by a trusted source, you can also enable trust all (*).

ConsumerConfig 类

@Configuration
@EnableKafka
public class KafkaConfig {

private ConsumerFactory<String, User> userConsumerFactory()
    {
        Map<String,Object> config=new ConcurrentHashMap<>();

        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"127.0.0.1:9092");
        config.put(ConsumerConfig.GROUP_ID_CONFIG,"group_json");
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, JsonDeserializer.class);
        config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS,JsonDeserializer.class.getClass());
        config.put(JsonDeserializer.TRUSTED_PACKAGES,"*");

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

 @Bean
    public ConcurrentKafkaListenerContainerFactory<String,User> userKafkaListenerContainerFactory()
    {
        ConcurrentKafkaListenerContainerFactory<String,User> factory
                =new ConcurrentKafkaListenerContainerFactory<>();

        factory.setConsumerFactory(userConsumerFactory());

        return factory;

    }
}

型号

@NoArgsConstructor
@AllArgsConstructor
@Getter
@Setter
public class User {

    private int userId;
    private String userName;
    private double salary;

    @Override
    public String toString() {
        return "User{" +
                "userId=" + userId +
                ", userName='" + userName + '\'' +
                ", salary=" + salary +
                '}';
    }
}

ProducerConfig 类

@Configuration
public class KafkaConfig {

    private ProducerFactory<String,User> producerFactory()
    {
        Map<String,Object> config=new ConcurrentHashMap<>();

          config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"127.0.0.1:9092");
          config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
          config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);

          return new DefaultKafkaProducerFactory<>(config);
    }

    @Bean
    public KafkaTemplate<String,User> kafkaTemplate()
    {
        return new KafkaTemplate<>(producerFactory());
    }


}

型号

@NoArgsConstructor
@AllArgsConstructor
@Getter
@Setter
public class User {

    private int userId;
    private String userName;
    private double salary;

    @Override
    public String toString() {
        return "User{" +
                "userId=" + userId +
                ", userName='" + userName + '\'' +
                ", salary=" + salary +
                '}';
    }
}

【问题讨论】:

    标签: java spring spring-boot apache-kafka spring-kafka


    【解决方案1】:

    您需要使用@Bean 对其进行注释并添加此配置:

        @Bean
        private ProducerFactory<String,User> producerFactory()
    {
        Map<String,Object> config=new ConcurrentHashMap<>();
    
          config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"127.0.0.1:9092");
          config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
          config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
    
          return new DefaultKafkaProducerFactory<>(config);
    }
    
    @Bean
        public Map<String, Object> producerConfigs() {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            return props;
        }
    
        @Bean
        public ProducerFactory<String, AccountEvent> producerFactory() {
            return new DefaultKafkaProducerFactory<>(producerConfigs());
        }
    
        @Bean
        public KafkaTemplate<String, AccountEvent> kafkaTemplate() {
            return new KafkaTemplate<>(producerFactory());
        }
    

    如果你想使用带有 kafka 基础设施的项目,我在我的 github 中有它。 https://github.com/gabryellr/banking-system

    【讨论】:

    • 不需要为这两种方法创建 bean 实例。顺便说一句,它没有工作
    【解决方案2】:

    这个

    config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
    

    应该是

    config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    

    但是,不清楚为什么您的受信任包配置没有被应用;建议你在JsonDeserializer.configure()方法中设置断点。

    编辑

    哦...

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

    当你像这样传入一个反序列化器实例时,属性不会被使用;您必须完全自己构建和配置反序列化器。

    既然你想使用ErrorHandlingDeserializer,你应该在这里传递一个。

    或者,将其更改为

    return new DefaultKafkaConsumerFactory<>(config);
    

    【讨论】:

    • 查看我的答案的编辑;你。正在覆盖构造函数中的反序列化器属性。
    • 我改了但是没用。当我将kafka-console-producer 用作生产者时,消费者应用程序会正确使用它,但是当我将生产者用作弹簧启动应用程序并发送时,就会发生此错误
    • 这毫无意义,
    【解决方案3】:

    我已经找到了发生此错误的原因。您可以看到错误部分模型类是在两个不同的包中创建的。在Producer中,模型类创建在com.benz.kafka.api.model包下,而在Consumer部分,模型类创建在com.benz.kafka.consumer.api.model包下。这是根本原因,我将com.benz.kafka.consumer.api.model 更改为com.benz.kafka.api.model 然后它就可以工作了。

    【讨论】:

      猜你喜欢
      • 2022-01-13
      • 1970-01-01
      • 1970-01-01
      • 2018-07-17
      • 2020-07-10
      • 2023-04-11
      • 2019-09-01
      • 1970-01-01
      • 2023-03-13
      相关资源
      最近更新 更多