From 964a193c4639d6c8f1881575de78a9053d7d52d8 Mon Sep 17 00:00:00 2001 From: Kirk True Date: Fri, 29 Mar 2024 14:43:56 -0700 Subject: [PATCH] =?UTF-8?q?KAFKA-16438:=20Update=20consumer=5Ftest.py?= =?UTF-8?q?=E2=80=99s=20static=20tests=20to=20support=20KIP-848=E2=80=99s?= =?UTF-8?q?=20group=20protocol=20config?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Migrated the following tests for the new consumer: - test_fencing_static_consumer - test_static_consumer_bounce - test_static_consumer_persisted_after_rejoin --- tests/kafkatest/tests/client/consumer_test.py | 29 ++++++++++--------- 1 file changed, 16 insertions(+), 13 deletions(-) diff --git a/tests/kafkatest/tests/client/consumer_test.py b/tests/kafkatest/tests/client/consumer_test.py index d4d6af9f2aad0..89c58e5f973e6 100644 --- a/tests/kafkatest/tests/client/consumer_test.py +++ b/tests/kafkatest/tests/client/consumer_test.py @@ -194,7 +194,7 @@ def test_consumer_bounce(self, clean_shutdown, bounce_mode, metadata_quorum=quor static_membership=[True, False], bounce_mode=["all", "rolling"], num_bounces=[5], - metadata_quorum=[quorum.zk], + metadata_quorum=[quorum.zk, quorum.isolated_kraft], use_new_coordinator=[False] ) @matrix( @@ -203,9 +203,10 @@ def test_consumer_bounce(self, clean_shutdown, bounce_mode, metadata_quorum=quor bounce_mode=["all", "rolling"], num_bounces=[5], metadata_quorum=[quorum.isolated_kraft], - use_new_coordinator=[True, False] + use_new_coordinator=[True], + group_protocol=consumer_group.all_group_protocols ) - def test_static_consumer_bounce(self, clean_shutdown, static_membership, bounce_mode, num_bounces, metadata_quorum=quorum.zk, use_new_coordinator=False): + def test_static_consumer_bounce(self, clean_shutdown, static_membership, bounce_mode, num_bounces, metadata_quorum=quorum.zk, use_new_coordinator=False, group_protocol=None): """ Verify correct static consumer behavior when the consumers in the group are restarted. In order to make sure the behavior of static members are different from dynamic ones, we take both static and dynamic @@ -226,7 +227,7 @@ def test_static_consumer_bounce(self, clean_shutdown, static_membership, bounce_ self.await_produced_messages(producer) self.session_timeout_sec = 60 - consumer = self.setup_consumer(self.TOPIC, static_membership=static_membership) + consumer = self.setup_consumer(self.TOPIC, static_membership=static_membership, group_protocol=group_protocol) consumer.start() self.await_all_members(consumer) @@ -268,15 +269,16 @@ def test_static_consumer_bounce(self, clean_shutdown, static_membership, bounce_ @cluster(num_nodes=7) @matrix( bounce_mode=["all", "rolling"], - metadata_quorum=[quorum.zk], + metadata_quorum=[quorum.zk, quorum.isolated_kraft], use_new_coordinator=[False] ) @matrix( bounce_mode=["all", "rolling"], metadata_quorum=[quorum.isolated_kraft], - use_new_coordinator=[True, False] + use_new_coordinator=[True], + group_protocol=consumer_group.all_group_protocols ) - def test_static_consumer_persisted_after_rejoin(self, bounce_mode, metadata_quorum=quorum.zk, use_new_coordinator=False): + def test_static_consumer_persisted_after_rejoin(self, bounce_mode, metadata_quorum=quorum.zk, use_new_coordinator=False, group_protocol=None): """ Verify that the updated member.id(updated_member_id) caused by static member rejoin would be persisted. If not, after the brokers rolling bounce, the migrated group coordinator would load the stale persisted member.id and @@ -291,7 +293,7 @@ def test_static_consumer_persisted_after_rejoin(self, bounce_mode, metadata_quor producer.start() self.await_produced_messages(producer) self.session_timeout_sec = 60 - consumer = self.setup_consumer(self.TOPIC, static_membership=True) + consumer = self.setup_consumer(self.TOPIC, static_membership=True, group_protocol=group_protocol) consumer.start() self.await_all_members(consumer) @@ -309,16 +311,17 @@ def test_static_consumer_persisted_after_rejoin(self, bounce_mode, metadata_quor @matrix( num_conflict_consumers=[1, 2], fencing_stage=["stable", "all"], - metadata_quorum=[quorum.zk], + metadata_quorum=[quorum.zk, quorum.isolated_kraft], use_new_coordinator=[False] ) @matrix( num_conflict_consumers=[1, 2], fencing_stage=["stable", "all"], metadata_quorum=[quorum.isolated_kraft], - use_new_coordinator=[True, False] + use_new_coordinator=[True], + group_protocol=consumer_group.all_group_protocols ) - def test_fencing_static_consumer(self, num_conflict_consumers, fencing_stage, metadata_quorum=quorum.zk, use_new_coordinator=False): + def test_fencing_static_consumer(self, num_conflict_consumers, fencing_stage, metadata_quorum=quorum.zk, use_new_coordinator=False, group_protocol=None): """ Verify correct static consumer behavior when there are conflicting consumers with same group.instance.id. @@ -335,10 +338,10 @@ def test_fencing_static_consumer(self, num_conflict_consumers, fencing_stage, me self.await_produced_messages(producer) self.session_timeout_sec = 60 - consumer = self.setup_consumer(self.TOPIC, static_membership=True) + consumer = self.setup_consumer(self.TOPIC, static_membership=True, group_protocol=group_protocol) self.num_consumers = num_conflict_consumers - conflict_consumer = self.setup_consumer(self.TOPIC, static_membership=True) + conflict_consumer = self.setup_consumer(self.TOPIC, static_membership=True, group_protocol=group_protocol) # wait original set of consumer to stable stage before starting conflict members. if fencing_stage == "stable":