【问题标题】:Apache Kafka Error Handling at Consumer End消费者端的 Apache Kafka 错误处理
【发布时间】:2023-01-10 22:54:24
【问题描述】:

您好,我正在使用 Apache Kafka 来消费来自另一个应用程序的消息。当消息反序列化或转换出现问题时,我想处理错误场景。我正在使用 Avro 架构来接收对象。

我实施了以下

@Configuration
@Slf4j
public class ConsumerConfig {
  @Bean
  ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
      ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
      ConsumerFactory<Object, Object> kafkaConsumerFactory) {
    ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    configurer.configure(factory, kafkaConsumerFactory);
    factory.setErrorHandler(((exception, data) -> {           
      log.error("Error in process with Exception {} and the record is {}", exception, data);
    }));
    return factory;
  }
}

但是如果我传递不同对象类型的消息,上面的代码不会处理它。我试图传递一个字符串,它抛出了错误但没有进入 Error Hdnaler。

org.apache.kafka.common.errors.InvalidConfigurationException: Schema being registered is incompatible with an earlier schema for subject "taas.cacib.lscsad-dev.queue.wwfdbtemp.Avros-value" io.confluent.kafka.schemaregistry.rest.exceptions.RestIncompatibleSchemaException: Schema being registered is incompatible with an earlier schema for subject "taas.cacib.lscsad-dev.queue.wwfdbtemp.Avros-value"

【问题讨论】:

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


    【解决方案1】:

    创建 kafka 主题后,您是否随时更改了方案?

    如果是这种情况,您将不得不重置 Kafka 主题以获取新方案,如果您想更改主题的方案,则需要创建一个新主题。

    【讨论】:

      【解决方案2】:

      使用ErrorHandlingDeserializer

      https://docs.spring.io/spring-kafka/docs/current/reference/html/#error-handling-deserializer

      然后,序列化错误将直接发送到错误处理程序,默认情况下,错误处理程序将此类错误视为致命错误,不会重试,而是直接将它们发送给恢复程序(默认记录)。

      【讨论】:

        猜你喜欢
        • 2019-11-25
        • 1970-01-01
        • 2020-12-14
        • 2015-06-05
        • 1970-01-01
        • 2019-01-24
        • 2021-09-22
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多