diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/DedicatedMirrorIntegrationTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/DedicatedMirrorIntegrationTest.java index e2db7b3865e12..2c57a4396f1d8 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/DedicatedMirrorIntegrationTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/DedicatedMirrorIntegrationTest.java @@ -47,8 +47,8 @@ public class DedicatedMirrorIntegrationTest { private static final Logger log = LoggerFactory.getLogger(DedicatedMirrorIntegrationTest.class); - private static final int TOPIC_CREATION_TIMEOUT_MS = 30_000; - private static final int TOPIC_REPLICATION_TIMEOUT_MS = 30_000; + private static final int TOPIC_CREATION_TIMEOUT_MS = 120_000; + private static final int TOPIC_REPLICATION_TIMEOUT_MS = 120_000; private Map kafkaClusters; private Map mirrorMakers; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java index b867a8ab13ec3..dac8e7b9bb7f5 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java @@ -1246,8 +1246,7 @@ void fenceZombieSourceTasks(final String connName, final Callback callback private void doFenceZombieSourceTasks(String connName, Callback callback) { log.trace("Performing zombie fencing request for connector {}", connName); - if (checkRebalanceNeeded(callback)) - return; + // We don't have to check for a pending rebalance here if (!isLeader()) callback.onCompletion(new NotLeaderException("Only the leader may perform zombie fencing.", leaderUrl()), null);