【问题标题】:Does spring Kafka batch listener commits db transaction in batch mode and in case of failure is the complete transaction rolled back?Spring Kafka 批处理侦听器是否以批处理模式提交数据库事务,如果失败,是否会回滚整个事务?
【发布时间】:2020-06-23 22:05:31
【问题描述】:

我有一个简单的要求来读取 kafka 消息并存储在数据库中。我在批处理侦听器模式下使用 spring kafka。我已经浏览了 spring kafka 文档,但仍然不清楚在批处理侦听器模式下使用 spring kafka 时,它是否以批处理模式提交 db 事务,如果失败,是否会回滚整个事务?

如果失败,它会再次寻找相同的记录集吗?

我有以下配置,

props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfigProperties.getBootstrapservers());
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaConfigProperties.getConsumer().getGroupid());
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, kafkaConfigProperties.getConsumer().getOffset());
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,250);
props.put(ApplicationConstant.KAFKA_SCHEMA_URL_PROPERTY, kafkaConfigProperties.getSchemaregistry());
@Bean
    public ConcurrentKafkaListenerContainerFactory<String, GenericRecord> kafkaListenerContainerFactory(KafkaConfigProperties kafkaConfigProperties) {
        ConcurrentKafkaListenerContainerFactory<String, GenericRecord> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        
        factory.setConsumerFactory(consumerFactory(kafkaConfigProperties));
        factory.setConcurrency(2);
        factory.setBatchListener(true);
        ContainerProperties containerProperties = factory.getContainerProperties();
        containerProperties.setAckOnError(false);
        containerProperties.setAckMode(AckMode.BATCH);
        return factory;
    }

【问题讨论】:

    标签: java spring apache-kafka spring-kafka


    【解决方案1】:

    您需要添加SeekToCurrentBatchErrorHandlerRecoveringBatchErrorHander 才能重播批处理。这是 2.5 及更高版本的默认错误处理程序。

    the documentation

    【讨论】:

    • 谢谢@GaryRussell,我想确保所有记录都作为批处理事务提交到数据库中,或者在发生任何异常时回滚,这是使用 kafka 批处理侦听器时的默认行为还是我需要添加 kafka 事务?
    • 该用例不需要 Kafka 事务,只需要 DB 事务管理器和侦听器上的 @Transactional(或它调用的东西);侦听器将在数据库事务中运行并在正常退出时提交;如果抛出异常则回滚,SeekToCurrentBatchErrorHandler 将重播批处理。
    • 谢谢@GaryRussell,如果事务回滚,DefaultAfterRollbackProcessor 不会提供类似的行为,所以我需要配置 SeekToCurrentBatchErrorHandler 吗?另外我使用的是 spring-kafka 版本 2.1.12,它不允许配置最大重试次数或固定退避。默认行为是什么,是否可以使用 2.1.12 版本的 kafka 配置无限尝试行为?
    • 否; DARP 仅用于 Kafka 事务;侦听器容器对您的数据库事务一无所知。你需要一个 STCBH。使用这样的旧版本,它将无限期地重试,没有后退。您确实需要升级到更新的版本才能获得改进的功能。不再支持 2.1.x。
    猜你喜欢
    • 2021-08-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多