From 69e9892a8c56be0948b781dc9163adfc6ecdc99f Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Sun, 26 Mar 2023 08:06:43 +0530 Subject: [PATCH 01/10] KAFKA-12525: Ignoring Stale status statuses when reading from Status topic --- .../storage/KafkaStatusBackingStore.java | 9 +++++- .../storage/KafkaStatusBackingStoreTest.java | 31 +++++++++++++++++++ 2 files changed, 39 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 20f47db2262a4..7d747548c93c3 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -588,7 +588,14 @@ private void readTaskStatus(String key, byte[] value) { synchronized (this) { log.trace("Received task {} status update {}", id, status); CacheEntry entry = getOrAdd(id); - entry.put(status); + // We add the status only if it's safe to do so. This is applicable + // when there are race conditions during frequent rebalances leading to + // the RUNNING status record with the latest generation followed by a stale + // UNASSIGNED status record belonging to the same or a previous generation + // from a worker which couldn't read the RUNNING record of the latest generation. + if (entry.canWriteSafely(status)) { + entry.put(status); + } } } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index 71fe342ee281b..660460a9f7567 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -309,6 +309,37 @@ public void readTaskState() { assertEquals(status, store.get(TASK)); } + @Test + public void readTaskStateShouldIgnoreStaleStatuses() { + byte[] value = new byte[0]; + String otherWorkerId = "anotherhost:8083"; + + // This worker sends a RUNNING status in the most recent generation + Map firstStatusRead = new HashMap<>(); + firstStatusRead.put("worker_id", otherWorkerId); + firstStatusRead.put("state", "RUNNING"); + firstStatusRead.put("generation", 10L); + + // Another worker still ends up producing an UNASSIGNED status before it could + // read the newer RUNNING status from above belonging to an older generation. + Map secondStatusRead = new HashMap<>(); + secondStatusRead.put("worker_id", WORKER_ID); + secondStatusRead.put("state", "UNASSIGNED"); + secondStatusRead.put("generation", 9L); + + when(converter.toConnectData(STATUS_TOPIC, value)) + .thenReturn(new SchemaAndValue(null, firstStatusRead)) + .thenReturn(new SchemaAndValue(null, secondStatusRead)); + + store.read(consumerRecord(0, "status-task-conn-0", value)); + store.read(consumerRecord(0, "status-task-conn-0", value)); + + verify(kafkaBasedLog, never()).send(anyString(), any(), any(Callback.class)); + // The latest task status should reflect RUNNING status from the newer generation + TaskStatus status = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, otherWorkerId, 10); + assertEquals(status, store.get(TASK)); + } + @Test public void deleteConnectorState() { final byte[] value = new byte[0]; From 42ce67413882c246e3f2223cb086ef18c1a48021 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Sun, 26 Mar 2023 08:15:22 +0530 Subject: [PATCH 02/10] Logging ignored status update --- .../kafka/connect/storage/KafkaStatusBackingStore.java | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 7d747548c93c3..a1e77e88d6005 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -588,13 +588,10 @@ private void readTaskStatus(String key, byte[] value) { synchronized (this) { log.trace("Received task {} status update {}", id, status); CacheEntry entry = getOrAdd(id); - // We add the status only if it's safe to do so. This is applicable - // when there are race conditions during frequent rebalances leading to - // the RUNNING status record with the latest generation followed by a stale - // UNASSIGNED status record belonging to the same or a previous generation - // from a worker which couldn't read the RUNNING record of the latest generation. if (entry.canWriteSafely(status)) { entry.put(status); + } else { + log.trace("Ignoring stale status update {} for task {}", status, id); } } } From 24ee4ee6e18910c3477cf7fe22fe2738cc9472e2 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Mon, 3 Apr 2023 08:46:45 +0530 Subject: [PATCH 03/10] Adding condition for Connector update status and ensuring tests capture ignroing stale updates from other workers only --- .../storage/KafkaStatusBackingStore.java | 6 +- .../storage/KafkaStatusBackingStoreTest.java | 67 +++++++++++++++++-- 2 files changed, 67 insertions(+), 6 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index a1e77e88d6005..8d61a810bbed2 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -562,7 +562,11 @@ private void readConnectorStatus(String key, byte[] value) { synchronized (this) { log.trace("Received connector {} status update {}", connector, status); CacheEntry entry = getOrAdd(connector); - entry.put(status); + if (entry.canWriteSafely(status)) { + entry.put(status); + } else { + log.trace("Ignoring stale status update {} for connector {}", status, connector); + } } } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index 660460a9f7567..6cd036ecbbc9a 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -27,6 +27,7 @@ import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.data.Struct; +import org.apache.kafka.connect.runtime.AbstractStatus; import org.apache.kafka.connect.runtime.ConnectorStatus; import org.apache.kafka.connect.runtime.TaskStatus; import org.apache.kafka.connect.runtime.WorkerConfig; @@ -229,6 +230,42 @@ public void putConnectorStateShouldOverride() { firstStatusRead.put("state", "RUNNING"); firstStatusRead.put("generation", 1L); + Map secondStatusRead = new HashMap<>(); + secondStatusRead.put("worker_id", WORKER_ID); + secondStatusRead.put("state", "UNASSIGNED"); + secondStatusRead.put("generation", 2L); + + when(converter.toConnectData(STATUS_TOPIC, value)) + .thenReturn(new SchemaAndValue(null, firstStatusRead)) + .thenReturn(new SchemaAndValue(null, secondStatusRead)); + + when(converter.fromConnectData(eq(STATUS_TOPIC), any(Schema.class), any(Struct.class))) + .thenReturn(value); + + doAnswer(invocation -> { + ((Callback) invocation.getArgument(2)).onCompletion(null, null); + store.read(consumerRecord(1, "status-connector-conn", value)); + return null; + }).when(kafkaBasedLog).send(eq("status-connector-conn"), eq(value), any(Callback.class)); + + store.read(consumerRecord(0, "status-connector-conn", value)); + + ConnectorStatus status = new ConnectorStatus(CONNECTOR, ConnectorStatus.State.UNASSIGNED, WORKER_ID, 2); + store.put(status); + assertEquals(status, store.get(CONNECTOR)); + } + + @Test + public void putConnectorStateShouldNotOverrideStaleStatus() { + final byte[] value = new byte[0]; + String otherWorkerId = "anotherhost:8083"; + + // the persisted came from a different host and has a newer generation + Map firstStatusRead = new HashMap<>(); + firstStatusRead.put("worker_id", otherWorkerId); + firstStatusRead.put("state", "RUNNING"); + firstStatusRead.put("generation", 1L); + Map secondStatusRead = new HashMap<>(); secondStatusRead.put("worker_id", WORKER_ID); secondStatusRead.put("state", "UNASSIGNED"); @@ -251,7 +288,8 @@ public void putConnectorStateShouldOverride() { ConnectorStatus status = new ConnectorStatus(CONNECTOR, ConnectorStatus.State.UNASSIGNED, WORKER_ID, 0); store.put(status); - assertEquals(status, store.get(CONNECTOR)); + ConnectorStatus expectedStatus = new ConnectorStatus(CONNECTOR, ConnectorStatus.State.RUNNING, otherWorkerId, 1); + assertEquals(expectedStatus, store.get(CONNECTOR)); } @Test @@ -310,7 +348,7 @@ public void readTaskState() { } @Test - public void readTaskStateShouldIgnoreStaleStatuses() { + public void readTaskStateShouldIgnoreStaleStatusesFromOtherWorkers() { byte[] value = new byte[0]; String otherWorkerId = "anotherhost:8083"; @@ -327,17 +365,36 @@ public void readTaskStateShouldIgnoreStaleStatuses() { secondStatusRead.put("state", "UNASSIGNED"); secondStatusRead.put("generation", 9L); + Map thirdStatusRead = new HashMap<>(); + thirdStatusRead.put("worker_id", otherWorkerId); + thirdStatusRead.put("state", "RUNNING"); + thirdStatusRead.put("generation", 8L); + when(converter.toConnectData(STATUS_TOPIC, value)) .thenReturn(new SchemaAndValue(null, firstStatusRead)) - .thenReturn(new SchemaAndValue(null, secondStatusRead)); + .thenReturn(new SchemaAndValue(null, secondStatusRead)) + .thenReturn(new SchemaAndValue(null, thirdStatusRead)); + + when(converter.fromConnectData(eq(STATUS_TOPIC), any(Schema.class), any(Struct.class))) + .thenReturn(value); + + doAnswer(invocation -> { + ((Callback) invocation.getArgument(2)).onCompletion(null, null); + store.read(consumerRecord(2, "status-task-conn-0", value)); + return null; + }).when(kafkaBasedLog).send(eq("status-task-conn-0"), eq(value), any(Callback.class)); store.read(consumerRecord(0, "status-task-conn-0", value)); - store.read(consumerRecord(0, "status-task-conn-0", value)); + store.read(consumerRecord(1, "status-task-conn-0", value)); - verify(kafkaBasedLog, never()).send(anyString(), any(), any(Callback.class)); // The latest task status should reflect RUNNING status from the newer generation TaskStatus status = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, otherWorkerId, 10); assertEquals(status, store.get(TASK)); + + // This stale status is from the same worker. In this case, the status should get updated + TaskStatus staleStatus = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, otherWorkerId, 8); + store.put(staleStatus); + assertEquals(staleStatus, store.get(TASK)); } @Test From 786f50d204cb51203a70dc7298a33a5d5104c44e Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Mon, 3 Apr 2023 22:20:33 +0530 Subject: [PATCH 04/10] Fixing checkstyle --- .../kafka/connect/storage/KafkaStatusBackingStoreTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index 6cd036ecbbc9a..f4dc8e8c052c0 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -27,7 +27,6 @@ import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.data.Struct; -import org.apache.kafka.connect.runtime.AbstractStatus; import org.apache.kafka.connect.runtime.ConnectorStatus; import org.apache.kafka.connect.runtime.TaskStatus; import org.apache.kafka.connect.runtime.WorkerConfig; From f442bc58f85bc7e6ebb6993bfbe7cddcb63e4513 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Mon, 17 Apr 2023 22:11:20 +0530 Subject: [PATCH 05/10] Adding more dedicated checks for UNASSIGNED and RUNTIME states --- .../storage/KafkaStatusBackingStore.java | 24 +++++---- .../storage/KafkaStatusBackingStoreTest.java | 49 +++---------------- 2 files changed, 22 insertions(+), 51 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 8d61a810bbed2..5753255cb459e 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -315,6 +315,7 @@ private > void send(final String key, this.generation = status.generation(); if (safeWrite && !entry.canWriteSafely(status)) return; + sequence = entry.increment(); } @@ -562,11 +563,7 @@ private void readConnectorStatus(String key, byte[] value) { synchronized (this) { log.trace("Received connector {} status update {}", connector, status); CacheEntry entry = getOrAdd(connector); - if (entry.canWriteSafely(status)) { - entry.put(status); - } else { - log.trace("Ignoring stale status update {} for connector {}", status, connector); - } + entry.put(status); } } @@ -592,11 +589,20 @@ private void readTaskStatus(String key, byte[] value) { synchronized (this) { log.trace("Received task {} status update {}", id, status); CacheEntry entry = getOrAdd(id); - if (entry.canWriteSafely(status)) { - entry.put(status); - } else { - log.trace("Ignoring stale status update {} for task {}", status, id); + + // During frequent rebalances, there could be a race condition because of which + // an UNASSIGNED state of a prior generation can be sent by a worker despite a + // RUNNING status in a newer generation by another worker because the first worker + // couldn't read the newer RUNNING status. This can lead to an inaccurate status + // representation even though the task might be actually running. + if (status.state() == ConnectorStatus.State.UNASSIGNED + && entry.get().state() == ConnectorStatus.State.RUNNING + && entry.get().generation() > status.generation()) { + log.trace("Ignoring stale status {} in favour of more upto date status {}", status, entry.get()); + return; } + + entry.put(status); } } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index f4dc8e8c052c0..3abe0befec1ff 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -229,42 +229,6 @@ public void putConnectorStateShouldOverride() { firstStatusRead.put("state", "RUNNING"); firstStatusRead.put("generation", 1L); - Map secondStatusRead = new HashMap<>(); - secondStatusRead.put("worker_id", WORKER_ID); - secondStatusRead.put("state", "UNASSIGNED"); - secondStatusRead.put("generation", 2L); - - when(converter.toConnectData(STATUS_TOPIC, value)) - .thenReturn(new SchemaAndValue(null, firstStatusRead)) - .thenReturn(new SchemaAndValue(null, secondStatusRead)); - - when(converter.fromConnectData(eq(STATUS_TOPIC), any(Schema.class), any(Struct.class))) - .thenReturn(value); - - doAnswer(invocation -> { - ((Callback) invocation.getArgument(2)).onCompletion(null, null); - store.read(consumerRecord(1, "status-connector-conn", value)); - return null; - }).when(kafkaBasedLog).send(eq("status-connector-conn"), eq(value), any(Callback.class)); - - store.read(consumerRecord(0, "status-connector-conn", value)); - - ConnectorStatus status = new ConnectorStatus(CONNECTOR, ConnectorStatus.State.UNASSIGNED, WORKER_ID, 2); - store.put(status); - assertEquals(status, store.get(CONNECTOR)); - } - - @Test - public void putConnectorStateShouldNotOverrideStaleStatus() { - final byte[] value = new byte[0]; - String otherWorkerId = "anotherhost:8083"; - - // the persisted came from a different host and has a newer generation - Map firstStatusRead = new HashMap<>(); - firstStatusRead.put("worker_id", otherWorkerId); - firstStatusRead.put("state", "RUNNING"); - firstStatusRead.put("generation", 1L); - Map secondStatusRead = new HashMap<>(); secondStatusRead.put("worker_id", WORKER_ID); secondStatusRead.put("state", "UNASSIGNED"); @@ -287,8 +251,7 @@ public void putConnectorStateShouldNotOverrideStaleStatus() { ConnectorStatus status = new ConnectorStatus(CONNECTOR, ConnectorStatus.State.UNASSIGNED, WORKER_ID, 0); store.put(status); - ConnectorStatus expectedStatus = new ConnectorStatus(CONNECTOR, ConnectorStatus.State.RUNNING, otherWorkerId, 1); - assertEquals(expectedStatus, store.get(CONNECTOR)); + assertEquals(status, store.get(CONNECTOR)); } @Test @@ -350,6 +313,7 @@ public void readTaskState() { public void readTaskStateShouldIgnoreStaleStatusesFromOtherWorkers() { byte[] value = new byte[0]; String otherWorkerId = "anotherhost:8083"; + String yetAnotherWorkerId = "yetanotherhost:8083"; // This worker sends a RUNNING status in the most recent generation Map firstStatusRead = new HashMap<>(); @@ -365,9 +329,9 @@ public void readTaskStateShouldIgnoreStaleStatusesFromOtherWorkers() { secondStatusRead.put("generation", 9L); Map thirdStatusRead = new HashMap<>(); - thirdStatusRead.put("worker_id", otherWorkerId); + thirdStatusRead.put("worker_id", yetAnotherWorkerId); thirdStatusRead.put("state", "RUNNING"); - thirdStatusRead.put("generation", 8L); + thirdStatusRead.put("generation", 1L); when(converter.toConnectData(STATUS_TOPIC, value)) .thenReturn(new SchemaAndValue(null, firstStatusRead)) @@ -390,8 +354,9 @@ public void readTaskStateShouldIgnoreStaleStatusesFromOtherWorkers() { TaskStatus status = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, otherWorkerId, 10); assertEquals(status, store.get(TASK)); - // This stale status is from the same worker. In this case, the status should get updated - TaskStatus staleStatus = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, otherWorkerId, 8); + // This status is from the another worker not necessarily belonging to the above group. + // In this case, the status should get update irrespective of whatever status info was present before. + TaskStatus staleStatus = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, yetAnotherWorkerId, 1); store.put(staleStatus); assertEquals(staleStatus, store.get(TASK)); } From 46c19ea46823f788beadb0837ae47d1fc40544ab Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Fri, 12 May 2023 16:51:23 +0530 Subject: [PATCH 06/10] Adding a null check when the CachedEntry is null for fixing some ITs --- .../kafka/connect/storage/KafkaStatusBackingStore.java | 2 +- .../connect/storage/KafkaStatusBackingStoreTest.java | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 5753255cb459e..73b21ef1a45b6 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -315,7 +315,6 @@ private > void send(final String key, this.generation = status.generation(); if (safeWrite && !entry.canWriteSafely(status)) return; - sequence = entry.increment(); } @@ -596,6 +595,7 @@ private void readTaskStatus(String key, byte[] value) { // couldn't read the newer RUNNING status. This can lead to an inaccurate status // representation even though the task might be actually running. if (status.state() == ConnectorStatus.State.UNASSIGNED + && entry.get() != null && entry.get().state() == ConnectorStatus.State.RUNNING && entry.get().generation() > status.generation()) { log.trace("Ignoring stale status {} in favour of more upto date status {}", status, entry.get()); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index 3abe0befec1ff..8c5dd78027c33 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -355,10 +355,10 @@ public void readTaskStateShouldIgnoreStaleStatusesFromOtherWorkers() { assertEquals(status, store.get(TASK)); // This status is from the another worker not necessarily belonging to the above group. - // In this case, the status should get update irrespective of whatever status info was present before. - TaskStatus staleStatus = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, yetAnotherWorkerId, 1); - store.put(staleStatus); - assertEquals(staleStatus, store.get(TASK)); + // In this case, the status should get updated irrespective of whatever status info was present before. + TaskStatus latestStatus = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, yetAnotherWorkerId, 1); + store.put(latestStatus); + assertEquals(latestStatus, store.get(TASK)); } @Test From 80e15de59b826ac0662898721ea75a38ad58cbb3 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Sat, 20 May 2023 09:36:48 +0530 Subject: [PATCH 07/10] Replacing ConnectorStatus with TaskStatus --- .../kafka/connect/storage/KafkaStatusBackingStore.java | 5 +++-- .../kafka/connect/storage/KafkaStatusBackingStoreTest.java | 4 ++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 73b21ef1a45b6..e80524f1d2ea7 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -31,6 +31,7 @@ import org.apache.kafka.common.utils.ThreadUtils; import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.connector.Task; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.data.SchemaBuilder; @@ -594,9 +595,9 @@ private void readTaskStatus(String key, byte[] value) { // RUNNING status in a newer generation by another worker because the first worker // couldn't read the newer RUNNING status. This can lead to an inaccurate status // representation even though the task might be actually running. - if (status.state() == ConnectorStatus.State.UNASSIGNED + if (status.state() == TaskStatus.State.UNASSIGNED && entry.get() != null - && entry.get().state() == ConnectorStatus.State.RUNNING + && entry.get().state() == TaskStatus.State.RUNNING && entry.get().generation() > status.generation()) { log.trace("Ignoring stale status {} in favour of more upto date status {}", status, entry.get()); return; diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index 8c5dd78027c33..4cad728ae65f7 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -351,12 +351,12 @@ public void readTaskStateShouldIgnoreStaleStatusesFromOtherWorkers() { store.read(consumerRecord(1, "status-task-conn-0", value)); // The latest task status should reflect RUNNING status from the newer generation - TaskStatus status = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, otherWorkerId, 10); + TaskStatus status = new TaskStatus(TASK, TaskStatus.State.RUNNING, otherWorkerId, 10); assertEquals(status, store.get(TASK)); // This status is from the another worker not necessarily belonging to the above group. // In this case, the status should get updated irrespective of whatever status info was present before. - TaskStatus latestStatus = new TaskStatus(TASK, ConnectorStatus.State.RUNNING, yetAnotherWorkerId, 1); + TaskStatus latestStatus = new TaskStatus(TASK, TaskStatus.State.RUNNING, yetAnotherWorkerId, 1); store.put(latestStatus); assertEquals(latestStatus, store.get(TASK)); } From f3c4e8b0bfc893b5d1eb34bbe8e457e0d2aa721c Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Sun, 21 May 2023 16:20:05 +0530 Subject: [PATCH 08/10] Checkstyle fix plus making generation check gte for the stale generation --- .../kafka/connect/storage/KafkaStatusBackingStore.java | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index e80524f1d2ea7..a9d2a06d46a6c 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -31,7 +31,6 @@ import org.apache.kafka.common.utils.ThreadUtils; import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.connector.Task; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.data.SchemaBuilder; @@ -591,14 +590,14 @@ private void readTaskStatus(String key, byte[] value) { CacheEntry entry = getOrAdd(id); // During frequent rebalances, there could be a race condition because of which - // an UNASSIGNED state of a prior generation can be sent by a worker despite a - // RUNNING status in a newer generation by another worker because the first worker - // couldn't read the newer RUNNING status. This can lead to an inaccurate status + // an UNASSIGNED state of a prior or same generation can be sent by a worker despite a + // RUNNING status by another worker because the first worker + // couldn't read the latest RUNNING status. This can lead to an inaccurate status // representation even though the task might be actually running. if (status.state() == TaskStatus.State.UNASSIGNED && entry.get() != null && entry.get().state() == TaskStatus.State.RUNNING - && entry.get().generation() > status.generation()) { + && entry.get().generation() >= status.generation()) { log.trace("Ignoring stale status {} in favour of more upto date status {}", status, entry.get()); return; } From 6dbdc44e0315133de481baa4f143f18b7f0bcd80 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Thu, 29 Jun 2023 00:33:06 +0530 Subject: [PATCH 09/10] Updating comment to reflect missed UNASSIGNED status in case of generation resets --- .../kafka/connect/storage/KafkaStatusBackingStore.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index a9d2a06d46a6c..0cd3c8396089f 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -594,6 +594,11 @@ private void readTaskStatus(String key, byte[] value) { // RUNNING status by another worker because the first worker // couldn't read the latest RUNNING status. This can lead to an inaccurate status // representation even though the task might be actually running. + // Note that this could also mean that when a generation reset happens, and an + // UNASSIGNED status is sent, then it would be ignored if the current status is RUNNING + // at a higher generation. But since it will be followed by a RUNNING or a different + // status message(at a lower generation) soon after, the misrepresentation of the UNASSIGNED + // status would be short-lived in most cases. if (status.state() == TaskStatus.State.UNASSIGNED && entry.get() != null && entry.get().state() == TaskStatus.State.RUNNING From 706ff84a826ac522996bf13fa29246bc2d71d3d7 Mon Sep 17 00:00:00 2001 From: Sagar Rao Date: Thu, 29 Jun 2023 00:36:48 +0530 Subject: [PATCH 10/10] Updating comment --- .../apache/kafka/connect/storage/KafkaStatusBackingStore.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 0cd3c8396089f..d1580e3a7c1d7 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -597,7 +597,7 @@ private void readTaskStatus(String key, byte[] value) { // Note that this could also mean that when a generation reset happens, and an // UNASSIGNED status is sent, then it would be ignored if the current status is RUNNING // at a higher generation. But since it will be followed by a RUNNING or a different - // status message(at a lower generation) soon after, the misrepresentation of the UNASSIGNED + // status message(at the lower generation) soon after, the misrepresentation of the UNASSIGNED // status would be short-lived in most cases. if (status.state() == TaskStatus.State.UNASSIGNED && entry.get() != null