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 51258f655ecd0..9826296e0c266 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 @@ -1104,7 +1104,9 @@ public void stopConnector(final String connName, final Callback callback) writeTaskConfigs(connName, Collections.emptyList()); configBackingStore.putTargetState(connName, TargetState.STOPPED); // Force a read of the new target state for the connector - refreshConfigSnapshot(workerSyncTimeoutMs); + if (!refreshConfigSnapshot(workerSyncTimeoutMs)) { + log.warn("Failed to read to end of config topic after writing the STOPPED target state for connector {}", connName); + } callback.onCompletion(null, null); return null;