diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/RetryUtil.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/RetryUtil.java index 9463f6ab2e463..64fd97ca6b691 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/RetryUtil.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/RetryUtil.java @@ -88,10 +88,13 @@ public static T retryUntilTimeout(Callable callable, Supplier des lastError = e; } - long millisRemaining = Math.max(0, end - System.currentTimeMillis()); - if (millisRemaining < retryBackoffMs) { - // exit when the time remaining is less than retryBackoffMs - break; + if (retryBackoffMs > 0) { + long millisRemaining = Math.max(0, end - System.currentTimeMillis()); + if (millisRemaining < retryBackoffMs) { + // exit when the time remaining is less than retryBackoffMs + break; + } + Utils.sleep(retryBackoffMs); } Utils.sleep(retryBackoffMs); } while (System.currentTimeMillis() < end); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/util/RetryUtilTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/util/RetryUtilTest.java index 05f021288019a..58c101bec3b02 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/util/RetryUtilTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/util/RetryUtilTest.java @@ -36,6 +36,8 @@ @RunWith(PowerMockRunner.class) public class RetryUtilTest { + private static final Duration TIMEOUT = Duration.ofSeconds(10); + private Callable mockCallable; private final Supplier testMsg = () -> "Test"; @@ -69,7 +71,7 @@ public void retriesEventuallySucceed() throws Exception { .thenThrow(new TimeoutException()) .thenReturn("success"); - assertEquals("success", RetryUtil.retryUntilTimeout(mockCallable, testMsg, Duration.ofMillis(100), 1)); + assertEquals("success", RetryUtil.retryUntilTimeout(mockCallable, testMsg, TIMEOUT, 1)); Mockito.verify(mockCallable, Mockito.times(4)).call(); } @@ -83,7 +85,7 @@ public void failWithNonRetriableException() throws Exception { .thenThrow(new TimeoutException("timeout")) .thenThrow(new NullPointerException("Non retriable")); NullPointerException e = assertThrows(NullPointerException.class, - () -> RetryUtil.retryUntilTimeout(mockCallable, testMsg, Duration.ofMillis(100), 0)); + () -> RetryUtil.retryUntilTimeout(mockCallable, testMsg, TIMEOUT, 0)); assertEquals("Non retriable", e.getMessage()); Mockito.verify(mockCallable, Mockito.times(6)).call(); } @@ -114,7 +116,7 @@ public void testNoBackoffTimeAndSucceed() throws Exception { .thenThrow(new TimeoutException()) .thenReturn("success"); - assertEquals("success", RetryUtil.retryUntilTimeout(mockCallable, testMsg, Duration.ofMillis(50), 0)); + assertEquals("success", RetryUtil.retryUntilTimeout(mockCallable, testMsg, TIMEOUT, 0)); Mockito.verify(mockCallable, Mockito.times(4)).call(); } @@ -151,7 +153,7 @@ public void testInvalidRetryTimeout() throws Exception { Mockito.when(mockCallable.call()) .thenThrow(new TimeoutException("timeout")) .thenReturn("success"); - assertEquals("success", RetryUtil.retryUntilTimeout(mockCallable, testMsg, Duration.ofMillis(100), -1)); + assertEquals("success", RetryUtil.retryUntilTimeout(mockCallable, testMsg, TIMEOUT, -1)); Mockito.verify(mockCallable, Mockito.times(2)).call(); }