【问题标题】:How to use Retry Tempate (backOffPolicy) with AfterRollbackProcessor in Spring Kafka 2.2.9.RELEASE如何在 Spring Kafka 2.2.9.RELEASE 中将 Retry Tempate (backOffPolicy) 与 AfterRollbackProcessor 一起使用
【发布时间】:2023-03-03 13:50:01
【问题描述】:

我正在使用带有 Kafka 和 MySQL 的 spring boot 2.1.9,并且还实现了一个链式事务管理器。

我想设置 backOffPolicy 以便在特定时间后重试。在新的 spring Kafka 版本中是可能的,但由于一些其他依赖项,我无法升级 spring boot。

到目前为止,我正在使用 AfterRollbackProcessor 处理失败的消息,现在我想使用 Spring Kafka 2.2.9.RELEASE 使用 AfterRollbackProcessor 实现 backoffPolicy。有什么方法可以实现吗?

这里是接收器配置文件:

@Configuration
@EnableKafka
public class KafkaReceiverConfig {

    // Kafka Server Configuration
    @Value("${kafka.servers}")
    private String kafkaServers;

    // Group Identifier
    @Value("${kafka.groupId}")
    private String groupId;

    // Kafka Max Retry Attempts
    @Value("${kafka.retry.maxAttempts:5}")
    private Integer retryMaxAttempts;

    // Kafka Max Retry Interval
    @Value("${kafka.retry.interval:180000}")
    private Long retryInterval;

    // Kafka Concurrency
    @Value("${kafka.concurrency:10}")
    private Integer concurrency;

    // Kafka Concurrency
    @Value("${kafka.poll.timeout:300}")
    private Integer pollTimeout;

    // Kafka Consumer Offset
    @Value("${kafka.consumer.auto-offset-reset:earliest}")
    private String offset = "earliest";

    @Value("${kafka.max.records:100}")
    private Integer maxPollRecords;

    @Value("${kafka.max.poll.interval.time:500000}")
    private Integer maxPollIntervalMs;

    @Value("${kafka.max.session.timeout:60000}")
    private Integer sessionTimoutMs;

    // Logger
    private static final Logger log = LoggerFactory.getLogger(KafkaReceiverConfig.class);

    /**
     * String Kafka Listener Container Factor
     * 
     * @return @see {@link KafkaListenerContainerFactory}
     */
    @Bean
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory(
            ChainedKafkaTransactionManager<String, String> chainedTM, MessageProducer messageProducer) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(concurrency);
        factory.getContainerProperties().setPollTimeout(pollTimeout);
        factory.getContainerProperties().setAckMode(AckMode.RECORD);
        factory.getContainerProperties().setSyncCommits(true);
        // factory.setRetryTemplate(retryTemplate());
        factory.getContainerProperties().setAckOnError(false);
        factory.getContainerProperties().setTransactionManager(chainedTM);
        // factory.setStatefulRetry(true);
        AfterRollbackProcessor<String, String> afterRollbackProcessor = new DefaultAfterRollbackProcessor<>(
                (record, exception) -> {
                    log.warn("failed to process kafka message (retries are exausted). topic name:" + record.topic()
                            + " value:" + record.value());
                    messageProducer.saveFailedMessage(record, exception);
                }, retryMaxAttempts);

        factory.setAfterRollbackProcessor(afterRollbackProcessor);
        log.debug("Kafka Receiver Config kafkaListenerContainerFactory created");
        return factory;
    }

    /**
     * String Consumer Factory
     * 
     * @return @see {@link ConsumerFactory}
     */
    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        log.debug("Kafka Receiver Config consumerFactory created");
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    /**
     * Consumer Configurations
     * 
     * @return @see {@link Map}
     */
    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new ConcurrentHashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaServers);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecords);
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPollIntervalMs);
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionTimoutMs);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, offset);
        props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
        log.debug("Kafka Receiver Config consumerConfigs created");
        return props;
    }

}

【问题讨论】:

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


    【解决方案1】:

    您可以使用侦听器重试,但它必须是有状态的(您已将其注释掉)。否则,重试将在事务中执行,这通常不是您想要的。

    使用有状态重试,模板在退出后抛出异常;然后后回滚处理器将执行重新搜索,以便重新处理记录。

    如您所说,在 2.3 中,我们在后回滚处理器中添加了一个 BackOff,以便更轻松地在一个地方配置所有内容。

    【讨论】:

    • 如果我对 afterRollbackProcessor 使用有状态重试,事务回滚是立即发生还是在下一次重试之前发生?因为如果重试失败,我将在每次重试之前获得日志跟踪和“事务回滚”消息(而不是在每次重试失败之后)。
    • 捕获异常时应用回退策略;在退避策略确定的睡眠后,异常被抛出到容器(并且事务被回滚)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-28
    • 2020-09-10
    • 1970-01-01
    • 2021-07-25
    • 2019-12-30
    • 2015-08-12
    相关资源
    最近更新 更多