Polish "Group Kafka back-off properties"

See gh-41335
This commit is contained in:
Andy Wilkinson
2024-07-11 12:17:00 +01:00
parent 14c9893371
commit 870739955b
3 changed files with 42 additions and 25 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2012-2023 the original author or authors.
* Copyright 2012-2024 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.

View File

@@ -1548,28 +1548,6 @@ public class KafkaProperties {
*/
private int attempts = 3;
/**
* Canonical backoff period. Used as an initial value in the exponential case,
* and as a minimum value in the uniform case.
*/
private Duration delay = Duration.ofSeconds(1);
/**
* Multiplier to use for generating the next backoff delay.
*/
private double multiplier = 0.0;
/**
* Maximum wait between retries. If less than the delay then the default of 30
* seconds is applied.
*/
private Duration maxDelay = Duration.ZERO;
/**
* Whether to have the backoff delays.
*/
private boolean randomBackOff = false;
public boolean isEnabled() {
return this.enabled;
}
@@ -1620,8 +1598,7 @@ public class KafkaProperties {
getBackoff().setMaxDelay(maxDelay);
}
@DeprecatedConfigurationProperty(replacement = "spring.kafka.retry.topic.backoff.random",
since = "3.4.0")
@DeprecatedConfigurationProperty(replacement = "spring.kafka.retry.topic.backoff.random", since = "3.4.0")
@Deprecated(since = "3.4.0", forRemoval = true)
public boolean isRandomBackOff() {
return getBackoff().isRandom();

View File

@@ -442,6 +442,22 @@ class KafkaAutoConfigurationTests {
@Test
void retryTopicConfigurationWithExponentialBackOff() {
this.contextRunner.withPropertyValues("spring.application.name=my-test-app",
"spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", "spring.kafka.retry.topic.enabled=true",
"spring.kafka.retry.topic.attempts=5", "spring.kafka.retry.topic.backoff.delay=100ms",
"spring.kafka.retry.topic.backoff.multiplier=2", "spring.kafka.retry.topic.backoff.max-delay=300ms")
.run((context) -> {
RetryTopicConfiguration configuration = context.getBean(RetryTopicConfiguration.class);
assertThat(configuration.getDestinationTopicProperties()).hasSize(5)
.extracting(DestinationTopic.Properties::delay, DestinationTopic.Properties::suffix)
.containsExactly(tuple(0L, ""), tuple(100L, "-retry-0"), tuple(200L, "-retry-1"),
tuple(300L, "-retry-2"), tuple(0L, "-dlt"));
});
}
@Test
@Deprecated(since = "3.4.0", forRemoval = true)
void retryTopicConfigurationWithExponentialBackOffUsingDeprecatedProperties() {
this.contextRunner.withPropertyValues("spring.application.name=my-test-app",
"spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", "spring.kafka.retry.topic.enabled=true",
"spring.kafka.retry.topic.attempts=5", "spring.kafka.retry.topic.delay=100ms",
@@ -471,6 +487,18 @@ class KafkaAutoConfigurationTests {
@Test
void retryTopicConfigurationWithFixedBackOff() {
this.contextRunner.withPropertyValues("spring.application.name=my-test-app",
"spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", "spring.kafka.retry.topic.enabled=true",
"spring.kafka.retry.topic.attempts=4", "spring.kafka.retry.topic.backoff.delay=2s")
.run(assertRetryTopicConfiguration(
(configuration) -> assertThat(configuration.getDestinationTopicProperties()).hasSize(3)
.extracting(DestinationTopic.Properties::delay)
.containsExactly(0L, 2000L, 0L)));
}
@Test
@Deprecated(since = "3.4.0", forRemoval = true)
void retryTopicConfigurationWithFixedBackOffUsingDeprecatedProperties() {
this.contextRunner.withPropertyValues("spring.application.name=my-test-app",
"spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", "spring.kafka.retry.topic.enabled=true",
"spring.kafka.retry.topic.attempts=4", "spring.kafka.retry.topic.delay=2s")
@@ -482,6 +510,18 @@ class KafkaAutoConfigurationTests {
@Test
void retryTopicConfigurationWithNoBackOff() {
this.contextRunner.withPropertyValues("spring.application.name=my-test-app",
"spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", "spring.kafka.retry.topic.enabled=true",
"spring.kafka.retry.topic.attempts=4", "spring.kafka.retry.topic.backoff.delay=0")
.run(assertRetryTopicConfiguration(
(configuration) -> assertThat(configuration.getDestinationTopicProperties()).hasSize(3)
.extracting(DestinationTopic.Properties::delay)
.containsExactly(0L, 0L, 0L)));
}
@Test
@Deprecated(since = "3.4.0", forRemoval = true)
void retryTopicConfigurationWithNoBackOffUsingDeprecatedProperties() {
this.contextRunner.withPropertyValues("spring.application.name=my-test-app",
"spring.kafka.bootstrap-servers=localhost:9092,localhost:9093", "spring.kafka.retry.topic.enabled=true",
"spring.kafka.retry.topic.attempts=4", "spring.kafka.retry.topic.delay=0")