From e21be806a3b54f68e131e5702bbab80acad1f433 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Thu, 13 Dec 2018 10:59:15 -0500 Subject: [PATCH 1/5] MINOR: fixes for making test more stable --- .../tests/streams/base_streams_test.py | 1 + .../streams_broker_down_resilience_test.py | 40 ++++++++----------- 2 files changed, 18 insertions(+), 23 deletions(-) diff --git a/tests/kafkatest/tests/streams/base_streams_test.py b/tests/kafkatest/tests/streams/base_streams_test.py index 320d4b2068b55..6e005dd6cfead 100644 --- a/tests/kafkatest/tests/streams/base_streams_test.py +++ b/tests/kafkatest/tests/streams/base_streams_test.py @@ -45,6 +45,7 @@ def get_producer(self, topic, num_messages, repeating_keys=None): topic, max_messages=num_messages, acks=1, + throughput=1000, repeating_keys=repeating_keys) def assert_produce_consume(self, diff --git a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py index 3cbf71390c9ba..f9da446616bcf 100644 --- a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py +++ b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py @@ -27,7 +27,8 @@ class StreamsBrokerDownResilience(BaseStreamsTest): inputTopic = "streamsResilienceSource" outputTopic = "streamsResilienceSink" client_id = "streams-broker-resilience-verify-consumer" - num_messages = 5 + num_messages = 10000 + message = "processed[0-9]*messages" def __init__(self, test_context): super(StreamsBrokerDownResilience, self).__init__(test_context, @@ -48,8 +49,6 @@ def test_streams_resilient_to_broker_down(self): processor = StreamsBrokerDownResilienceService(self.test_context, self.kafka, self.get_configs()) processor.start() - # until KIP-91 is merged we'll only send 5 messages to assert Kafka Streams is running before taking the broker down - # After KIP-91 is merged we'll continue to send messages the duration of the test self.assert_produce_consume(self.inputTopic, self.outputTopic, self.client_id, @@ -103,14 +102,13 @@ def test_streams_runs_with_broker_down_initially(self): self.outputTopic, self.client_id, "running_with_broker_down_initially", - num_messages=9, + num_messages=self.num_messages, timeout_sec=120) - message = "processed3messages" # need to show all 3 instances processed messages - self.wait_for_verification(processor, message, processor.STDOUT_FILE) - self.wait_for_verification(processor_2, message, processor_2.STDOUT_FILE) - self.wait_for_verification(processor_3, message, processor_3.STDOUT_FILE) + self.wait_for_verification(processor, self.message, processor.STDOUT_FILE) + self.wait_for_verification(processor_2, self.message, processor_2.STDOUT_FILE) + self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) self.kafka.stop() @@ -136,14 +134,12 @@ def test_streams_should_scale_in_while_brokers_down(self): self.outputTopic, self.client_id, "waiting for rebalance to complete", - num_messages=9, + num_messages=self.num_messages, timeout_sec=120) - message = "processed3messages" - - self.wait_for_verification(processor, message, processor.STDOUT_FILE) - self.wait_for_verification(processor_2, message, processor_2.STDOUT_FILE) - self.wait_for_verification(processor_3, message, processor_3.STDOUT_FILE) + self.wait_for_verification(processor, self.message, processor.STDOUT_FILE) + self.wait_for_verification(processor_2, self.message, processor_2.STDOUT_FILE) + self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) node = self.kafka.leader(self.inputTopic) self.kafka.stop_node(node) @@ -161,10 +157,10 @@ def test_streams_should_scale_in_while_brokers_down(self): self.outputTopic, self.client_id, "sending_message_after_stopping_streams_instance_bouncing_broker", - num_messages=9, + num_messages=self.num_messages, timeout_sec=120) - self.wait_for_verification(processor_3, "processed9messages", processor_3.STDOUT_FILE) + self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) self.kafka.stop() @@ -190,14 +186,12 @@ def test_streams_should_failover_while_brokers_down(self): self.outputTopic, self.client_id, "waiting for rebalance to complete", - num_messages=9, + num_messages=self.num_messages, timeout_sec=120) - message = "processed3messages" - - self.wait_for_verification(processor, message, processor.STDOUT_FILE) - self.wait_for_verification(processor_2, message, processor_2.STDOUT_FILE) - self.wait_for_verification(processor_3, message, processor_3.STDOUT_FILE) + self.wait_for_verification(processor, self.message, processor.STDOUT_FILE) + self.wait_for_verification(processor_2, self.message, processor_2.STDOUT_FILE) + self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) node = self.kafka.leader(self.inputTopic) self.kafka.stop_node(node) @@ -212,7 +206,7 @@ def test_streams_should_failover_while_brokers_down(self): self.outputTopic, self.client_id, "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", - num_messages=9, + num_messages=self.num_messages, timeout_sec=120) self.kafka.stop() From 5e0af0550b9d8ab7b4a0c92805a18262df656a9a Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Fri, 14 Dec 2018 11:49:52 -0500 Subject: [PATCH 2/5] MINOR: Wait for streams to re-connect with broker --- .../tests/streams/streams_broker_down_resilience_test.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py index f9da446616bcf..bfce494ec5e4f 100644 --- a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py +++ b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py @@ -62,6 +62,13 @@ def test_streams_resilient_to_broker_down(self): self.kafka.start_node(node) + connected = 'Discovered group coordinator' + + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + monitor.wait_until(connected, + timeout_sec=120, + err_msg=("Never saw output '%s' on " % connected) + str(processor.node.account)) + self.assert_produce_consume(self.inputTopic, self.outputTopic, self.client_id, From c6a851e430468e916a2e39212e32fea5f3207200 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Sat, 15 Dec 2018 17:35:31 -0500 Subject: [PATCH 3/5] MINOR: Move restart of kafka node inside grabbing the monitor --- .../tests/streams/streams_broker_down_resilience_test.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py index bfce494ec5e4f..3ec3f71b0a592 100644 --- a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py +++ b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py @@ -60,13 +60,12 @@ def test_streams_resilient_to_broker_down(self): time.sleep(broker_down_time_in_seconds) - self.kafka.start_node(node) - connected = 'Discovered group coordinator' with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + self.kafka.start_node(node) monitor.wait_until(connected, - timeout_sec=120, + timeout_sec=180, err_msg=("Never saw output '%s' on " % connected) + str(processor.node.account)) self.assert_produce_consume(self.inputTopic, From c9d542e16dfb274c0b3921615932a5ccc38c2680 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Sun, 16 Dec 2018 17:38:15 -0500 Subject: [PATCH 4/5] MINOR: More clean up and close up other timing error gaps --- .../streams_broker_down_resilience_test.py | 212 +++++++++++++----- 1 file changed, 152 insertions(+), 60 deletions(-) diff --git a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py index 3ec3f71b0a592..ffc92fa851233 100644 --- a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py +++ b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py @@ -29,6 +29,7 @@ class StreamsBrokerDownResilience(BaseStreamsTest): client_id = "streams-broker-resilience-verify-consumer" num_messages = 10000 message = "processed[0-9]*messages" + connected = "Discovered group coordinator" def __init__(self, test_context): super(StreamsBrokerDownResilience, self).__init__(test_context, @@ -60,13 +61,11 @@ def test_streams_resilient_to_broker_down(self): time.sleep(broker_down_time_in_seconds) - connected = 'Discovered group coordinator' - with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: self.kafka.start_node(node) - monitor.wait_until(connected, - timeout_sec=180, - err_msg=("Never saw output '%s' on " % connected) + str(processor.node.account)) + monitor.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw output '%s' on " % self.connected) + str(processor.node.account)) self.assert_produce_consume(self.inputTopic, self.outputTopic, @@ -100,21 +99,45 @@ def test_streams_runs_with_broker_down_initially(self): self.wait_for_verification(processor_2, broker_unavailable_message, processor_2.LOG_FILE, 10) self.wait_for_verification(processor_3, broker_unavailable_message, processor_3.LOG_FILE, 10) - # now start broker - self.kafka.start_node(node) - - # assert streams can process when starting with broker down - self.assert_produce_consume(self.inputTopic, - self.outputTopic, - self.client_id, - "running_with_broker_down_initially", - num_messages=self.num_messages, - timeout_sec=120) - - # need to show all 3 instances processed messages - self.wait_for_verification(processor, self.message, processor.STDOUT_FILE) - self.wait_for_verification(processor_2, self.message, processor_2.STDOUT_FILE) - self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor_1: + with processor_2.node.account.monitor_log(processor_2.LOG_FILE) as monitor_2: + with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor_3: + self.kafka.start_node(node) + + monitor_1.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor.node.account)) + monitor_2.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor_2.node.account)) + monitor_3.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor_3.node.account)) + + with processor.node.account.monitor_log(processor.STDOUT_FILE) as monitor_1: + with processor_2.node.account.monitor_log(processor_2.STDOUT_FILE) as monitor_2: + with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: + + self.assert_produce(self.inputTopic, + "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + num_messages=self.num_messages, + timeout_sec=120) + + monitor_1.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor.node.account)) + monitor_2.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_2.node.account)) + monitor_3.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) + + self.assert_consume(self.client_id, + "sending_message_after_stopping_streams_instance_bouncing_broker", + self.outputTopic, + num_messages=self.num_messages, + timeout_sec=120) self.kafka.stop() @@ -130,22 +153,40 @@ def test_streams_should_scale_in_while_brokers_down(self): processor_2.start() processor_3 = StreamsBrokerDownResilienceService(self.test_context, self.kafka, configs) - processor_3.start() # need to wait for rebalance once - self.wait_for_verification(processor_3, "State transition from REBALANCING to RUNNING", processor_3.LOG_FILE) - - # assert streams can process when starting with broker up - self.assert_produce_consume(self.inputTopic, - self.outputTopic, - self.client_id, - "waiting for rebalance to complete", - num_messages=self.num_messages, - timeout_sec=120) - - self.wait_for_verification(processor, self.message, processor.STDOUT_FILE) - self.wait_for_verification(processor_2, self.message, processor_2.STDOUT_FILE) - self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) + rebalance = "State transition from REBALANCING to RUNNING" + with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor: + processor_3.start() + + monitor.wait_until(rebalance, + timeout_sec=120, + err_msg=("Never saw output '%s' on " % rebalance) + str(processor_3.node.account)) + + with processor.node.account.monitor_log(processor.STDOUT_FILE) as monitor_1: + with processor_2.node.account.monitor_log(processor_2.STDOUT_FILE) as monitor_2: + with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: + + self.assert_produce(self.inputTopic, + "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + num_messages=self.num_messages, + timeout_sec=120) + + monitor_1.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor.node.account)) + monitor_2.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_2.node.account)) + monitor_3.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) + + self.assert_consume(self.client_id, + "sending_message_after_stopping_streams_instance_bouncing_broker", + self.outputTopic, + num_messages=self.num_messages, + timeout_sec=120) node = self.kafka.leader(self.inputTopic) self.kafka.stop_node(node) @@ -157,7 +198,12 @@ def test_streams_should_scale_in_while_brokers_down(self): self.wait_for_verification(processor, shutdown_message, processor.STDOUT_FILE) self.wait_for_verification(processor_2, shutdown_message, processor_2.STDOUT_FILE) - self.kafka.start_node(node) + with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor_3: + self.kafka.start_node(node) + + monitor_3.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor_3.node.account)) self.assert_produce_consume(self.inputTopic, self.outputTopic, @@ -166,8 +212,6 @@ def test_streams_should_scale_in_while_brokers_down(self): num_messages=self.num_messages, timeout_sec=120) - self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) - self.kafka.stop() def test_streams_should_failover_while_brokers_down(self): @@ -182,22 +226,40 @@ def test_streams_should_failover_while_brokers_down(self): processor_2.start() processor_3 = StreamsBrokerDownResilienceService(self.test_context, self.kafka, configs) - processor_3.start() # need to wait for rebalance once - self.wait_for_verification(processor_3, "State transition from REBALANCING to RUNNING", processor_3.LOG_FILE) - - # assert streams can process when starting with broker up - self.assert_produce_consume(self.inputTopic, - self.outputTopic, - self.client_id, - "waiting for rebalance to complete", - num_messages=self.num_messages, - timeout_sec=120) - - self.wait_for_verification(processor, self.message, processor.STDOUT_FILE) - self.wait_for_verification(processor_2, self.message, processor_2.STDOUT_FILE) - self.wait_for_verification(processor_3, self.message, processor_3.STDOUT_FILE) + rebalance = "State transition from REBALANCING to RUNNING" + with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor: + processor_3.start() + + monitor.wait_until(rebalance, + timeout_sec=120, + err_msg=("Never saw output '%s' on " % rebalance) + str(processor_3.node.account)) + + with processor.node.account.monitor_log(processor.STDOUT_FILE) as monitor_1: + with processor_2.node.account.monitor_log(processor_2.STDOUT_FILE) as monitor_2: + with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: + + self.assert_produce(self.inputTopic, + "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + num_messages=self.num_messages, + timeout_sec=120) + + monitor_1.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor.node.account)) + monitor_2.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_2.node.account)) + monitor_3.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) + + self.assert_consume(self.client_id, + "sending_message_after_stopping_streams_instance_bouncing_broker", + self.outputTopic, + num_messages=self.num_messages, + timeout_sec=120) node = self.kafka.leader(self.inputTopic) self.kafka.stop_node(node) @@ -206,13 +268,43 @@ def test_streams_should_failover_while_brokers_down(self): processor_2.abortThenRestart() processor_3.abortThenRestart() - self.kafka.start_node(node) - - self.assert_produce_consume(self.inputTopic, - self.outputTopic, - self.client_id, - "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", - num_messages=self.num_messages, - timeout_sec=120) - + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor_1: + with processor_2.node.account.monitor_log(processor_2.LOG_FILE) as monitor_2: + with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor_3: + self.kafka.start_node(node) + + monitor_1.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor.node.account)) + monitor_2.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor_2.node.account)) + monitor_3.wait_until(self.connected, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.connected) + str(processor_3.node.account)) + + with processor.node.account.monitor_log(processor.STDOUT_FILE) as monitor_1: + with processor_2.node.account.monitor_log(processor_2.STDOUT_FILE) as monitor_2: + with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: + + self.assert_produce(self.inputTopic, + "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + num_messages=self.num_messages, + timeout_sec=120) + + monitor_1.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor.node.account)) + monitor_2.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_2.node.account)) + monitor_3.wait_until(self.message, + timeout_sec=120, + err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) + + self.assert_consume(self.client_id, + "sending_message_after_stopping_streams_instance_bouncing_broker", + self.outputTopic, + num_messages=self.num_messages, + timeout_sec=120) self.kafka.stop() From b7cd8aadbe12bbeaa540a70836c74b10968bf57e Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Sun, 16 Dec 2018 17:46:34 -0500 Subject: [PATCH 5/5] MINOR: Clean up messages --- .../streams_broker_down_resilience_test.py | 48 +++++++++---------- 1 file changed, 24 insertions(+), 24 deletions(-) diff --git a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py index ffc92fa851233..ee5feaea1486a 100644 --- a/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py +++ b/tests/kafkatest/tests/streams/streams_broker_down_resilience_test.py @@ -29,7 +29,7 @@ class StreamsBrokerDownResilience(BaseStreamsTest): client_id = "streams-broker-resilience-verify-consumer" num_messages = 10000 message = "processed[0-9]*messages" - connected = "Discovered group coordinator" + connected_message = "Discovered group coordinator" def __init__(self, test_context): super(StreamsBrokerDownResilience, self).__init__(test_context, @@ -63,9 +63,9 @@ def test_streams_resilient_to_broker_down(self): with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: self.kafka.start_node(node) - monitor.wait_until(self.connected, + monitor.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw output '%s' on " % self.connected) + str(processor.node.account)) + err_msg=("Never saw output '%s' on " % self.connected_message) + str(processor.node.account)) self.assert_produce_consume(self.inputTopic, self.outputTopic, @@ -104,22 +104,22 @@ def test_streams_runs_with_broker_down_initially(self): with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor_3: self.kafka.start_node(node) - monitor_1.wait_until(self.connected, + monitor_1.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor.node.account)) - monitor_2.wait_until(self.connected, + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor.node.account)) + monitor_2.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor_2.node.account)) - monitor_3.wait_until(self.connected, + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor_2.node.account)) + monitor_3.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor_3.node.account)) + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor_3.node.account)) with processor.node.account.monitor_log(processor.STDOUT_FILE) as monitor_1: with processor_2.node.account.monitor_log(processor_2.STDOUT_FILE) as monitor_2: with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: self.assert_produce(self.inputTopic, - "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + "sending_message_after_broker_down_initially", num_messages=self.num_messages, timeout_sec=120) @@ -134,7 +134,7 @@ def test_streams_runs_with_broker_down_initially(self): err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) self.assert_consume(self.client_id, - "sending_message_after_stopping_streams_instance_bouncing_broker", + "consuming_message_after_broker_down_initially", self.outputTopic, num_messages=self.num_messages, timeout_sec=120) @@ -168,7 +168,7 @@ def test_streams_should_scale_in_while_brokers_down(self): with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: self.assert_produce(self.inputTopic, - "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + "sending_message_normal_broker_start", num_messages=self.num_messages, timeout_sec=120) @@ -183,7 +183,7 @@ def test_streams_should_scale_in_while_brokers_down(self): err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) self.assert_consume(self.client_id, - "sending_message_after_stopping_streams_instance_bouncing_broker", + "consuming_message_normal_broker_start", self.outputTopic, num_messages=self.num_messages, timeout_sec=120) @@ -201,9 +201,9 @@ def test_streams_should_scale_in_while_brokers_down(self): with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor_3: self.kafka.start_node(node) - monitor_3.wait_until(self.connected, + monitor_3.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor_3.node.account)) + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor_3.node.account)) self.assert_produce_consume(self.inputTopic, self.outputTopic, @@ -241,7 +241,7 @@ def test_streams_should_failover_while_brokers_down(self): with processor_3.node.account.monitor_log(processor_3.STDOUT_FILE) as monitor_3: self.assert_produce(self.inputTopic, - "sending_message_after_hard_bouncing_streams_instance_bouncing_broker", + "sending_message_after_normal_broker_start", num_messages=self.num_messages, timeout_sec=120) @@ -256,7 +256,7 @@ def test_streams_should_failover_while_brokers_down(self): err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) self.assert_consume(self.client_id, - "sending_message_after_stopping_streams_instance_bouncing_broker", + "consuming_message_after_normal_broker_start", self.outputTopic, num_messages=self.num_messages, timeout_sec=120) @@ -273,15 +273,15 @@ def test_streams_should_failover_while_brokers_down(self): with processor_3.node.account.monitor_log(processor_3.LOG_FILE) as monitor_3: self.kafka.start_node(node) - monitor_1.wait_until(self.connected, + monitor_1.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor.node.account)) - monitor_2.wait_until(self.connected, + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor.node.account)) + monitor_2.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor_2.node.account)) - monitor_3.wait_until(self.connected, + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor_2.node.account)) + monitor_3.wait_until(self.connected_message, timeout_sec=120, - err_msg=("Never saw '%s' on " % self.connected) + str(processor_3.node.account)) + err_msg=("Never saw '%s' on " % self.connected_message) + str(processor_3.node.account)) with processor.node.account.monitor_log(processor.STDOUT_FILE) as monitor_1: with processor_2.node.account.monitor_log(processor_2.STDOUT_FILE) as monitor_2: @@ -303,7 +303,7 @@ def test_streams_should_failover_while_brokers_down(self): err_msg=("Never saw '%s' on " % self.message) + str(processor_3.node.account)) self.assert_consume(self.client_id, - "sending_message_after_stopping_streams_instance_bouncing_broker", + "consuming_message_after_stopping_streams_instance_bouncing_broker", self.outputTopic, num_messages=self.num_messages, timeout_sec=120)