From de1c09348a155f48bec8b28a2936c38bd4d6f54c Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Tue, 3 Dec 2019 11:27:48 -0800 Subject: [PATCH 01/12] KAFKA-9184: Add a lightweight method to confirm connection to the broker coordinator is up from WorkerGroupMember --- .../runtime/distributed/WorkerGroupMember.java | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java index 4819db5da7417..6dab4549058b3 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java @@ -151,10 +151,23 @@ public void stop() { stop(false); } + /** + * Ensure that the connection to the broker coordinator is up and that the worker is an + * active member of the group. + */ public void ensureActive() { coordinator.poll(0); } + /** + * A lightweight wrapper call that answers whether there's an active connection to the broker + * coordinator. + * @return true if there is no active connection to a broker coordinator + */ + public boolean coordinatorUnknown() { + return coordinator.coordinatorUnknown(); + } + public void poll(long timeout) { if (timeout < 0) throw new IllegalArgumentException("Timeout must not be negative"); From ce4f95e98ed312716b684b75aebc1472a1370dea Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Tue, 3 Dec 2019 13:04:18 -0800 Subject: [PATCH 02/12] KAFKA-9184: Check connectivity with broker coordinator in intervals and stop tasks if coordinator is unreachable --- .../distributed/WorkerCoordinator.java | 20 ++++++++++++++++++- 1 file changed, 19 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java index a8ddbcc955b44..12e47e7d1d75c 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java @@ -40,6 +40,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.TimeUnit; import static org.apache.kafka.common.message.JoinGroupRequestData.JoinGroupRequestProtocolCollection; import static org.apache.kafka.common.message.JoinGroupResponseData.JoinGroupResponseMember; @@ -66,6 +67,7 @@ public class WorkerCoordinator extends AbstractCoordinator implements Closeable private volatile ConnectProtocolCompatibility currentConnectProtocol; private final ConnectAssignor eagerAssignor; private final ConnectAssignor incrementalAssignor; + private final int coordinatorDiscoveryTimeoutMs; /** * Initialize the coordination manager. @@ -98,6 +100,7 @@ public WorkerCoordinator(GroupRebalanceConfig config, this.incrementalAssignor = new IncrementalCooperativeAssignor(logContext, time, maxDelay); this.eagerAssignor = new EagerAssignor(logContext); this.currentConnectProtocol = protocolCompatibility; + this.coordinatorDiscoveryTimeoutMs = config.heartbeatIntervalMs; } @Override @@ -124,7 +127,22 @@ public void poll(long timeout) { do { if (coordinatorUnknown()) { - ensureCoordinatorReady(time.timer(Long.MAX_VALUE)); + log.debug("Broker coordinator is marked unknown. Attempting discovery with a timeout of {}ms", + coordinatorDiscoveryTimeoutMs); + if (ensureCoordinatorReady(time.timer(coordinatorDiscoveryTimeoutMs))) { + log.debug("Broker coordinator is ready"); + } else { + log.debug("Can not connect to broker coordinator"); + if (assignmentSnapshot != null && !assignmentSnapshot.failed()) { + log.info("Broker coordinator was unreachable for {}ms. Revoking previous assignment {} to " + + "avoid running tasks while not being a member the group", coordinatorDiscoveryTimeoutMs, assignmentSnapshot); + listener.onRevoked(assignmentSnapshot.leader(), assignmentSnapshot.connectors(), assignmentSnapshot.tasks()); + assignmentSnapshot.connectors().clear(); + assignmentSnapshot.tasks().clear(); + assignmentSnapshot.revokedConnectors().clear(); + assignmentSnapshot.revokedTasks().clear(); + } + } now = time.milliseconds(); } From 9b37c176f558bff3fc9edd336a209035ba5be210 Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Tue, 3 Dec 2019 13:15:48 -0800 Subject: [PATCH 03/12] KAFKA-9184: Reset rebalance delay when there are no lost tasks --- .../IncrementalCooperativeAssignor.java | 23 ++++++++++++++++--- .../distributed/WorkerCoordinator.java | 10 ++++++++ 2 files changed, 30 insertions(+), 3 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java index 91e1f7c8a73e9..6d1afbfa93b96 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java @@ -153,6 +153,9 @@ private Long ensureLeaderConfig(long maxOffset, WorkerCoordinator coordinator) { protected Map performTaskAssignment(String leaderId, long maxOffset, Map memberConfigs, WorkerCoordinator coordinator, short protocolVersion) { + log.debug("Performing task assignment during generation: {} with memberId: {}", + coordinator.generationId(), coordinator.memberId()); + // Base set: The previous assignment of connectors-and-tasks is a standalone snapshot that // can be used to calculate derived sets log.debug("Previous assignments: {}", previousAssignment); @@ -350,6 +353,7 @@ protected void handleLostAssignments(ConnectorsAndTasks lostAssignments, ConnectorsAndTasks newSubmissions, List completeWorkerAssignment) { if (lostAssignments.isEmpty()) { + resetDelay(); return; } @@ -359,6 +363,7 @@ protected void handleLostAssignments(ConnectorsAndTasks lostAssignments, if (scheduledRebalance > 0 && now >= scheduledRebalance) { // delayed rebalance expired and it's time to assign resources + log.debug("Delayed rebalance expired. Reassigning lost tasks"); Optional candidateWorkerLoad = Optional.empty(); if (!candidateWorkersForReassignment.isEmpty()) { candidateWorkerLoad = pickCandidateWorkerForReassignment(completeWorkerAssignment); @@ -366,15 +371,15 @@ protected void handleLostAssignments(ConnectorsAndTasks lostAssignments, if (candidateWorkerLoad.isPresent()) { WorkerLoad workerLoad = candidateWorkerLoad.get(); + log.debug("A candidate worker has been found to assign lost tasks: {}", workerLoad.worker()); lostAssignments.connectors().forEach(workerLoad::assign); lostAssignments.tasks().forEach(workerLoad::assign); } else { + log.debug("No single candidate worker was found to assign lost tasks. Treating lost tasks as new tasks"); newSubmissions.connectors().addAll(lostAssignments.connectors()); newSubmissions.tasks().addAll(lostAssignments.tasks()); } - candidateWorkersForReassignment.clear(); - scheduledRebalance = 0; - delay = 0; + resetDelay(); } else { candidateWorkersForReassignment .addAll(candidateWorkersForReassignment(completeWorkerAssignment)); @@ -382,17 +387,29 @@ protected void handleLostAssignments(ConnectorsAndTasks lostAssignments, // a delayed rebalance is in progress, but it's not yet time to reassign // unaccounted resources delay = calculateDelay(now); + log.debug("Delayed rebalance in progress. Task reassignment is postponed. New computed rebalance delay: {}", delay); } else { // This means scheduledRebalance == 0 // We could also also extract the current minimum delay from the group, to make // independent of consecutive leader failures, but this optimization is skipped // at the moment delay = maxDelay; + log.debug("Resetting rebalance delay to the max: {}. scheduledRebalance: {} now: {} diff scheduledRebalance - now: {}", + delay, scheduledRebalance, now, scheduledRebalance - now); } scheduledRebalance = now + delay; } } + private void resetDelay(){ + candidateWorkersForReassignment.clear(); + scheduledRebalance = 0; + if (delay != 0) { + log.debug("Resetting delay from previous value: {} to 0", delay); + } + delay = 0; + } + private Set candidateWorkersForReassignment(List completeWorkerAssignment) { return completeWorkerAssignment.stream() .filter(WorkerLoad::isEmpty) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java index 12e47e7d1d75c..e038cbd9f3b03 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java @@ -245,6 +245,16 @@ public String memberId() { return JoinGroupRequest.UNKNOWN_MEMBER_ID; } + /** + * Return the current generation. The generation refers to this worker's knowledge with + * respect to which generation is the latest one and, therefore, this information is local. + * + * @return the generation ID or -1 if no generation is defined + */ + public int generationId() { + return super.generation().generationId; + } + private boolean isLeader() { return assignmentSnapshot != null && memberId().equals(assignmentSnapshot.leader()); } From 1005a6151f4de726121fc8422a174ec0aa0837ae Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Tue, 3 Dec 2019 21:40:15 -0800 Subject: [PATCH 04/12] KAFKA-9184: Remove unused import --- .../kafka/connect/runtime/distributed/WorkerCoordinator.java | 1 - 1 file changed, 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java index e038cbd9f3b03..091a92112c298 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java @@ -40,7 +40,6 @@ import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.concurrent.TimeUnit; import static org.apache.kafka.common.message.JoinGroupRequestData.JoinGroupRequestProtocolCollection; import static org.apache.kafka.common.message.JoinGroupResponseData.JoinGroupResponseMember; From 9560d0bd29861152320805fbc30952062f7cefb6 Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Tue, 3 Dec 2019 22:30:07 -0800 Subject: [PATCH 05/12] KAFKA-9184: Set assignmentSnapshot to null after stopping tasks --- .../kafka/connect/runtime/distributed/WorkerCoordinator.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java index 091a92112c298..e5d841a2cf418 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java @@ -136,10 +136,7 @@ public void poll(long timeout) { log.info("Broker coordinator was unreachable for {}ms. Revoking previous assignment {} to " + "avoid running tasks while not being a member the group", coordinatorDiscoveryTimeoutMs, assignmentSnapshot); listener.onRevoked(assignmentSnapshot.leader(), assignmentSnapshot.connectors(), assignmentSnapshot.tasks()); - assignmentSnapshot.connectors().clear(); - assignmentSnapshot.tasks().clear(); - assignmentSnapshot.revokedConnectors().clear(); - assignmentSnapshot.revokedTasks().clear(); + assignmentSnapshot = null; } } now = time.milliseconds(); From 01caf295ead17a64c04966c4f81e84ba32b0614d Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Wed, 4 Dec 2019 00:24:34 -0800 Subject: [PATCH 06/12] KAFKA-9184: Extend embedded kafka cluster to allow restarts and add integration test --- .../ConnectWorkerIntegrationTest.java | 71 +++++++++++++++++-- .../util/clusters/EmbeddedKafkaCluster.java | 50 ++++++++++--- 2 files changed, 103 insertions(+), 18 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java index 7cdfa7dcb9df0..5be2c45cfc3a0 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.connect.integration; import org.apache.kafka.connect.runtime.AbstractStatus; +import org.apache.kafka.connect.runtime.distributed.DistributedConfig; import org.apache.kafka.connect.runtime.rest.entities.ConnectorStateInfo; import org.apache.kafka.connect.storage.StringConverter; import org.apache.kafka.connect.util.clusters.EmbeddedConnectCluster; @@ -65,29 +66,27 @@ public class ConnectWorkerIntegrationTest { private static final int NUM_WORKERS = 3; private static final String CONNECTOR_NAME = "simple-source"; + private EmbeddedConnectCluster.Builder connectBuilder; private EmbeddedConnectCluster connect; + Map workerProps = new HashMap<>(); + Properties brokerProps = new Properties(); @Before public void setup() throws IOException { // setup Connect worker properties - Map workerProps = new HashMap<>(); workerProps.put(OFFSET_COMMIT_INTERVAL_MS_CONFIG, String.valueOf(OFFSET_COMMIT_INTERVAL_MS)); workerProps.put(CONNECTOR_CLIENT_POLICY_CLASS_CONFIG, "All"); // setup Kafka broker properties - Properties brokerProps = new Properties(); brokerProps.put("auto.create.topics.enable", String.valueOf(false)); // build a Connect cluster backed by Kafka and Zk - connect = new EmbeddedConnectCluster.Builder() + connectBuilder = new EmbeddedConnectCluster.Builder() .name("connect-cluster") .numWorkers(NUM_WORKERS) .workerProps(workerProps) .brokerProps(brokerProps) - .build(); - - // start the clusters - connect.start(); + .maskExitProcedures(true); // true is the default, setting here as example } @After @@ -102,6 +101,10 @@ public void close() { */ @Test public void testAddAndRemoveWorker() throws Exception { + connect = connectBuilder.build(); + // start the clusters + connect.start(); + int numTasks = 4; // create test topic connect.kafka().createTopic("test-topic", NUM_TOPIC_PARTITIONS); @@ -149,6 +152,10 @@ public void testAddAndRemoveWorker() throws Exception { */ @Test public void testRestartFailedTask() throws Exception { + connect = connectBuilder.build(); + // start the clusters + connect.start(); + int numTasks = 1; // Properties for the source connector. The task should fail at startup due to the bad broker address. @@ -180,6 +187,56 @@ public void testRestartFailedTask() throws Exception { CONNECTOR_SETUP_DURATION_MS, "Connector tasks are not all in running state."); } + /** + * Verify that a set of tasks restarts correctly after a broker goes offline and back online + */ + @Test + public void testBrokerCoordinator() throws Exception { + workerProps.put(DistributedConfig.SCHEDULED_REBALANCE_MAX_DELAY_MS_CONFIG, String.valueOf(5000)); + connect = connectBuilder.workerProps(workerProps).build(); + // start the clusters + connect.start(); + int numTasks = 4; + // create test topic + connect.kafka().createTopic("test-topic", NUM_TOPIC_PARTITIONS); + + // setup up props for the sink connector + Map props = new HashMap<>(); + props.put(CONNECTOR_CLASS_CONFIG, MonitorableSourceConnector.class.getSimpleName()); + props.put(TASKS_MAX_CONFIG, String.valueOf(numTasks)); + props.put("topic", "test-topic"); + props.put("throughput", String.valueOf(1)); + props.put("messages.per.poll", String.valueOf(10)); + props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + + waitForCondition(() -> assertWorkersUp(NUM_WORKERS).orElse(false), + WORKER_SETUP_DURATION_MS, "Initial group of workers did not start in time."); + + // start a source connector + connect.configureConnector(CONNECTOR_NAME, props); + + waitForCondition(() -> assertConnectorAndTasksRunning(CONNECTOR_NAME, numTasks).orElse(false), + CONNECTOR_SETUP_DURATION_MS, "Connector tasks did not start in time."); + + connect.kafka().stopOnlyKafka(); + + waitForCondition(() -> assertWorkersUp(NUM_WORKERS).orElse(false), + WORKER_SETUP_DURATION_MS, "Group of workers did not remain the same after broker shutdown"); + + connect.kafka().startOnlyKafkaOnSamePorts(); + + // Allow for the workers to discover that the coordinator is unavailable + Thread.sleep(TimeUnit.SECONDS.toMillis(10)); + + waitForCondition(() -> assertWorkersUp(NUM_WORKERS).orElse(false), + WORKER_SETUP_DURATION_MS, "Group of workers did not remain the same within the " + + "designated time."); + + waitForCondition(() -> assertConnectorAndTasksRunning(CONNECTOR_NAME, numTasks).orElse(false), + CONNECTOR_SETUP_DURATION_MS, "Connector tasks did not start in time."); + } + /** * Confirm that the requested number of workers is up and running. * diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java b/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java index 948d54b57bd51..d36ae5b078e18 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedKafkaCluster.java @@ -81,6 +81,8 @@ public class EmbeddedKafkaCluster extends ExternalResource { private final KafkaServer[] brokers; private final Properties brokerConfig; private final Time time = new MockTime(); + private final int[] currentBrokerPorts; + private final String[] currentBrokerLogDirs; private EmbeddedZookeeper zookeeper = null; private ListenerName listenerName = new ListenerName("PLAINTEXT"); @@ -89,6 +91,8 @@ public class EmbeddedKafkaCluster extends ExternalResource { public EmbeddedKafkaCluster(final int numBrokers, final Properties brokerConfig) { brokers = new KafkaServer[numBrokers]; + currentBrokerPorts = new int[numBrokers]; + currentBrokerLogDirs = new String[numBrokers]; this.brokerConfig = brokerConfig; } @@ -102,11 +106,20 @@ protected void after() { stop(); } + public void startOnlyKafkaOnSamePorts() throws IOException { + start(currentBrokerPorts, currentBrokerLogDirs); + } + private void start() throws IOException { + // pick a random port zookeeper = new EmbeddedZookeeper(); + Arrays.fill(currentBrokerPorts, 0); + Arrays.fill(currentBrokerLogDirs, null); + start(currentBrokerPorts, currentBrokerLogDirs); + } + private void start(int[] brokerPorts, String[] logDirs) throws IOException { brokerConfig.put(KafkaConfig$.MODULE$.ZkConnectProp(), zKConnectString()); - brokerConfig.put(KafkaConfig$.MODULE$.PortProp(), 0); // pick a random port putIfAbsent(brokerConfig, KafkaConfig$.MODULE$.HostNameProp(), "localhost"); putIfAbsent(brokerConfig, KafkaConfig$.MODULE$.DeleteTopicEnableProp(), true); @@ -121,8 +134,11 @@ private void start() throws IOException { for (int i = 0; i < brokers.length; i++) { brokerConfig.put(KafkaConfig$.MODULE$.BrokerIdProp(), i); - brokerConfig.put(KafkaConfig$.MODULE$.LogDirProp(), createLogDir()); + currentBrokerLogDirs[i] = logDirs[i] == null ? createLogDir() : currentBrokerLogDirs[i]; + brokerConfig.put(KafkaConfig$.MODULE$.LogDirProp(), currentBrokerLogDirs[i]); + brokerConfig.put(KafkaConfig$.MODULE$.PortProp(), brokerPorts[i]); brokers[i] = TestUtils.createServer(new KafkaConfig(brokerConfig, true), time); + currentBrokerPorts[i] = brokers[i].boundPort(listenerName); } Map producerProps = new HashMap<>(); @@ -132,8 +148,15 @@ private void start() throws IOException { producer = new KafkaProducer<>(producerProps); } + public void stopOnlyKafka() { + stop(false, false); + } + private void stop() { + stop(true, true); + } + private void stop(boolean deleteLogDirs, boolean stopZK) { try { producer.close(); } catch (Exception e) { @@ -151,19 +174,24 @@ private void stop() { } } - for (KafkaServer broker : brokers) { - try { - log.info("Cleaning up kafka log dirs at {}", broker.config().logDirs()); - CoreUtils.delete(broker.config().logDirs()); - } catch (Throwable t) { - String msg = String.format("Could not clean up log dirs for broker at %s", address(broker)); - log.error(msg, t); - throw new RuntimeException(msg, t); + if (deleteLogDirs) { + for (KafkaServer broker : brokers) { + try { + log.info("Cleaning up kafka log dirs at {}", broker.config().logDirs()); + CoreUtils.delete(broker.config().logDirs()); + } catch (Throwable t) { + String msg = String.format("Could not clean up log dirs for broker at %s", + address(broker)); + log.error(msg, t); + throw new RuntimeException(msg, t); + } } } try { - zookeeper.shutdown(); + if (stopZK) { + zookeeper.shutdown(); + } } catch (Throwable t) { String msg = String.format("Could not shutdown zookeeper at %s", zKConnectString()); log.error(msg, t); From 6fdae3af9de45a85526d1b3d8814cb1a5916ed89 Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Wed, 4 Dec 2019 00:47:32 -0800 Subject: [PATCH 07/12] KAFKA-9184: Fix checkstyle --- .../runtime/distributed/IncrementalCooperativeAssignor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java index 6d1afbfa93b96..c6d7cc62e1145 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java @@ -401,7 +401,7 @@ protected void handleLostAssignments(ConnectorsAndTasks lostAssignments, } } - private void resetDelay(){ + private void resetDelay() { candidateWorkersForReassignment.clear(); scheduledRebalance = 0; if (delay != 0) { From 6e6eb0a94c2959b77c8623e2f84d98f3911022f4 Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Wed, 4 Dec 2019 01:00:27 -0800 Subject: [PATCH 08/12] KAFKA-9184: Remove unused method --- .../connect/runtime/distributed/WorkerGroupMember.java | 9 --------- 1 file changed, 9 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java index 6dab4549058b3..cc052c2907726 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java @@ -159,15 +159,6 @@ public void ensureActive() { coordinator.poll(0); } - /** - * A lightweight wrapper call that answers whether there's an active connection to the broker - * coordinator. - * @return true if there is no active connection to a broker coordinator - */ - public boolean coordinatorUnknown() { - return coordinator.coordinatorUnknown(); - } - public void poll(long timeout) { if (timeout < 0) throw new IllegalArgumentException("Timeout must not be negative"); From 4e9a6167d022e9757d7fdf247e5f4ca345b6487a Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Wed, 4 Dec 2019 02:09:42 -0800 Subject: [PATCH 09/12] KAFKA-9184: Adapt unit test to the additional debug calls --- .../IncrementalCooperativeAssignorTest.java | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java index c46d59bf8058a..a5e4ef0ca1bea 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java @@ -170,6 +170,8 @@ public void testTaskAssignmentWhenWorkerJoins() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -234,6 +236,8 @@ public void testTaskAssignmentWhenWorkerLeavesPermanently() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -314,6 +318,8 @@ public void testTaskAssignmentWhenWorkerBounces() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -369,6 +375,8 @@ public void testTaskAssignmentWhenLeaderLeavesPermanently() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -438,6 +446,8 @@ public void testTaskAssignmentWhenLeaderBounces() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -507,6 +517,8 @@ public void testTaskAssignmentWhenFirstAssignmentAttemptFails() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -563,6 +575,8 @@ public void testTaskAssignmentWhenSubsequentAssignmentAttemptFails() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test @@ -596,6 +610,8 @@ public void testTaskAssignmentWhenConnectorsAreDeleted() { verify(coordinator, times(rebalanceNum)).configSnapshot(); verify(coordinator, times(rebalanceNum)).leaderState(any()); + verify(coordinator, times(rebalanceNum)).generationId(); + verify(coordinator, times(rebalanceNum)).memberId(); } @Test From f7ded4f187efe4af067189ea68e539ab658dd14a Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Wed, 4 Dec 2019 02:11:09 -0800 Subject: [PATCH 10/12] KAFKA-9184: Add more specific logs to DistributedHerder --- .../connect/runtime/distributed/DistributedHerder.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) 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 eb618fb0053f3..f3861ddd59bf9 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 @@ -1538,10 +1538,11 @@ public void onAssigned(ExtendedAssignment assignment, int generation) { short priorProtocolVersion = currentProtocolVersion; DistributedHerder.this.currentProtocolVersion = member.currentProtocolVersion(); log.info( - "Joined group at generation {} with protocol version {} and got assignment: {}", + "Joined group at generation {} with protocol version {} and got assignment: {} with rebalance delay: {}", generation, DistributedHerder.this.currentProtocolVersion, - assignment + assignment, + assignment.delay() ); synchronized (DistributedHerder.this) { DistributedHerder.this.assignment = assignment; @@ -1611,12 +1612,13 @@ public void onRevoked(String leader, Collection connectors, Collection Date: Wed, 4 Dec 2019 09:39:26 -0800 Subject: [PATCH 11/12] KAFKA-9184: Declare assignmentSnapshot volatile and use local copies of the reference --- .../distributed/WorkerCoordinator.java | 52 +++++++++++-------- 1 file changed, 30 insertions(+), 22 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java index e5d841a2cf418..71f351df9ea63 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java @@ -56,7 +56,7 @@ public class WorkerCoordinator extends AbstractCoordinator implements Closeable private final Logger log; private final String restUrl; private final ConfigBackingStore configStorage; - private ExtendedAssignment assignmentSnapshot; + private volatile ExtendedAssignment assignmentSnapshot; private ClusterConfigState configSnapshot; private final WorkerRebalanceListener listener; private final ConnectProtocolCompatibility protocolCompatibility; @@ -132,10 +132,11 @@ public void poll(long timeout) { log.debug("Broker coordinator is ready"); } else { log.debug("Can not connect to broker coordinator"); - if (assignmentSnapshot != null && !assignmentSnapshot.failed()) { + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + if (localAssignmentSnapshot != null && !localAssignmentSnapshot.failed()) { log.info("Broker coordinator was unreachable for {}ms. Revoking previous assignment {} to " + - "avoid running tasks while not being a member the group", coordinatorDiscoveryTimeoutMs, assignmentSnapshot); - listener.onRevoked(assignmentSnapshot.leader(), assignmentSnapshot.connectors(), assignmentSnapshot.tasks()); + "avoid running tasks while not being a member the group", coordinatorDiscoveryTimeoutMs, localAssignmentSnapshot); + listener.onRevoked(localAssignmentSnapshot.leader(), localAssignmentSnapshot.connectors(), localAssignmentSnapshot.tasks()); assignmentSnapshot = null; } } @@ -166,7 +167,8 @@ public void poll(long timeout) { @Override public JoinGroupRequestProtocolCollection metadata() { configSnapshot = configStorage.snapshot(); - ExtendedWorkerState workerState = new ExtendedWorkerState(restUrl, configSnapshot.offset(), assignmentSnapshot); + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + ExtendedWorkerState workerState = new ExtendedWorkerState(restUrl, configSnapshot.offset(), localAssignmentSnapshot); switch (protocolCompatibility) { case EAGER: return ConnectProtocol.metadataRequest(workerState); @@ -194,17 +196,18 @@ protected void onJoinComplete(int generation, String memberId, String protocol, listener.onRevoked(newAssignment.leader(), newAssignment.revokedConnectors(), newAssignment.revokedTasks()); } - if (assignmentSnapshot != null) { - assignmentSnapshot.connectors().removeAll(newAssignment.revokedConnectors()); - assignmentSnapshot.tasks().removeAll(newAssignment.revokedTasks()); - log.debug("After revocations snapshot of assignment: {}", assignmentSnapshot); - newAssignment.connectors().addAll(assignmentSnapshot.connectors()); - newAssignment.tasks().addAll(assignmentSnapshot.tasks()); + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + if (localAssignmentSnapshot != null) { + localAssignmentSnapshot.connectors().removeAll(newAssignment.revokedConnectors()); + localAssignmentSnapshot.tasks().removeAll(newAssignment.revokedTasks()); + log.debug("After revocations snapshot of assignment: {}", localAssignmentSnapshot); + newAssignment.connectors().addAll(localAssignmentSnapshot.connectors()); + newAssignment.tasks().addAll(localAssignmentSnapshot.tasks()); } log.debug("Augmented new assignment: {}", newAssignment); } assignmentSnapshot = newAssignment; - listener.onAssigned(assignmentSnapshot, generation); + listener.onAssigned(newAssignment, generation); } @Override @@ -218,19 +221,21 @@ protected Map performAssignment(String leaderId, String prot protected void onJoinPrepare(int generation, String memberId) { log.info("Rebalance started"); leaderState(null); + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; if (currentConnectProtocol == EAGER) { - log.debug("Revoking previous assignment {}", assignmentSnapshot); - if (assignmentSnapshot != null && !assignmentSnapshot.failed()) - listener.onRevoked(assignmentSnapshot.leader(), assignmentSnapshot.connectors(), assignmentSnapshot.tasks()); + log.debug("Revoking previous assignment {}", localAssignmentSnapshot); + if (localAssignmentSnapshot != null && !localAssignmentSnapshot.failed()) + listener.onRevoked(localAssignmentSnapshot.leader(), localAssignmentSnapshot.connectors(), localAssignmentSnapshot.tasks()); } else { log.debug("Cooperative rebalance triggered. Keeping assignment {} until it's " - + "explicitly revoked.", assignmentSnapshot); + + "explicitly revoked.", localAssignmentSnapshot); } } @Override protected boolean rejoinNeededOrPending() { - return super.rejoinNeededOrPending() || (assignmentSnapshot == null || assignmentSnapshot.failed()) || rejoinRequested; + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + return super.rejoinNeededOrPending() || (localAssignmentSnapshot == null || localAssignmentSnapshot.failed()) || rejoinRequested; } @Override @@ -252,7 +257,8 @@ public int generationId() { } private boolean isLeader() { - return assignmentSnapshot != null && memberId().equals(assignmentSnapshot.leader()); + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + return localAssignmentSnapshot != null && memberId().equals(localAssignmentSnapshot.leader()); } public String ownerUrl(String connector) { @@ -330,20 +336,22 @@ public WorkerCoordinatorMetrics(Metrics metrics, String metricGrpPrefix) { Measurable numConnectors = new Measurable() { @Override public double measure(MetricConfig config, long now) { - if (assignmentSnapshot == null) { + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + if (localAssignmentSnapshot == null) { return 0.0; } - return assignmentSnapshot.connectors().size(); + return localAssignmentSnapshot.connectors().size(); } }; Measurable numTasks = new Measurable() { @Override public double measure(MetricConfig config, long now) { - if (assignmentSnapshot == null) { + final ExtendedAssignment localAssignmentSnapshot = assignmentSnapshot; + if (localAssignmentSnapshot == null) { return 0.0; } - return assignmentSnapshot.tasks().size(); + return localAssignmentSnapshot.tasks().size(); } }; From a031894a59ce13839b6e5d54bf22152fc5640f88 Mon Sep 17 00:00:00 2001 From: Konstantine Karantasis Date: Wed, 4 Dec 2019 10:06:16 -0800 Subject: [PATCH 12/12] KAFKA-9184: Add some more time between phases in the integration test --- .../integration/ConnectWorkerIntegrationTest.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java index 5be2c45cfc3a0..4a704664b69db 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectWorkerIntegrationTest.java @@ -224,15 +224,22 @@ public void testBrokerCoordinator() throws Exception { waitForCondition(() -> assertWorkersUp(NUM_WORKERS).orElse(false), WORKER_SETUP_DURATION_MS, "Group of workers did not remain the same after broker shutdown"); + // Allow for the workers to discover that the coordinator is unavailable, wait is + // heartbeat timeout * 2 + 4sec + Thread.sleep(TimeUnit.SECONDS.toMillis(10)); + connect.kafka().startOnlyKafkaOnSamePorts(); - // Allow for the workers to discover that the coordinator is unavailable + // Allow for the kafka brokers to come back online Thread.sleep(TimeUnit.SECONDS.toMillis(10)); waitForCondition(() -> assertWorkersUp(NUM_WORKERS).orElse(false), WORKER_SETUP_DURATION_MS, "Group of workers did not remain the same within the " + "designated time."); + // Allow for the workers to rebalance and reach a steady state + Thread.sleep(TimeUnit.SECONDS.toMillis(10)); + waitForCondition(() -> assertConnectorAndTasksRunning(CONNECTOR_NAME, numTasks).orElse(false), CONNECTOR_SETUP_DURATION_MS, "Connector tasks did not start in time."); }