From 9bc633dd61bd35bd5cf2f5e38e7ad63b5599f16a Mon Sep 17 00:00:00 2001 From: Yash Mayya Date: Tue, 23 May 2023 22:04:20 +0530 Subject: [PATCH 1/2] MINOR: Handle the config topic read timeout edge case in DistributedHerder's stopConnector method --- .../kafka/connect/runtime/distributed/DistributedHerder.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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..38414051c6bd7 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)) { + throw new ConnectException("Failed to read to end of config topic"); + } callback.onCompletion(null, null); return null; From bd1df7b2206e71f4ec00a84f10baede0643e3d02 Mon Sep 17 00:00:00 2001 From: Yash Mayya Date: Sun, 4 Jun 2023 23:22:55 +0530 Subject: [PATCH 2/2] Log a warning message instead of throwing an exception when timing out on the config topic read after writing a stopped target state --- .../kafka/connect/runtime/distributed/DistributedHerder.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 38414051c6bd7..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 @@ -1105,7 +1105,7 @@ public void stopConnector(final String connName, final Callback callback) configBackingStore.putTargetState(connName, TargetState.STOPPED); // Force a read of the new target state for the connector if (!refreshConfigSnapshot(workerSyncTimeoutMs)) { - throw new ConnectException("Failed to read to end of config topic"); + log.warn("Failed to read to end of config topic after writing the STOPPED target state for connector {}", connName); } callback.onCompletion(null, null);