Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 19 additions & 10 deletions tests/kafkatest/tests/core/kraft_upgrade_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,12 @@
from ducktape.mark.resource import cluster
from ducktape.utils.util import wait_until
from kafkatest.services.console_consumer import ConsoleConsumer
from kafkatest.services.kafka import KafkaService
from kafkatest.services.kafka import KafkaService, consumer_group
from kafkatest.services.kafka.quorum import isolated_kraft, combined_kraft
from kafkatest.services.verifiable_producer import VerifiableProducer
from kafkatest.tests.produce_consume_validate import ProduceConsumeValidateTest
from kafkatest.utils import is_int
from kafkatest.version import LATEST_3_1, LATEST_3_2, LATEST_3_3, LATEST_3_4, LATEST_3_5, \
from kafkatest.version import LATEST_3_1, LATEST_3_2, LATEST_3_3, LATEST_3_4, LATEST_3_5, LATEST_3_7, \
DEV_BRANCH, KafkaVersion, LATEST_STABLE_METADATA_VERSION

#
Expand Down Expand Up @@ -74,7 +74,7 @@ def perform_version_change(self, from_kafka_version):
self.logger.info("Changing metadata.version to %s" % LATEST_STABLE_METADATA_VERSION)
self.kafka.upgrade_metadata_version(LATEST_STABLE_METADATA_VERSION)

def run_upgrade(self, from_kafka_version):
def run_upgrade(self, from_kafka_version, group_protocol):
"""Test upgrade of Kafka broker cluster from various versions to the current version

from_kafka_version is a Kafka version to upgrade from.
Expand All @@ -101,7 +101,8 @@ def run_upgrade(self, from_kafka_version):
version=KafkaVersion(from_kafka_version))
self.consumer = ConsoleConsumer(self.test_context, self.num_consumers, self.kafka,
self.topic, new_consumer=True, consumer_timeout_ms=30000,
message_validator=is_int, version=KafkaVersion(from_kafka_version))
message_validator=is_int, version=KafkaVersion(from_kafka_version),
consumer_properties=consumer_group.maybe_set_group_protocol(group_protocol))
self.run_produce_consume_validate(core_test_action=lambda: self.perform_version_change(from_kafka_version))
cluster_id = self.kafka.cluster_id()
assert cluster_id is not None
Expand All @@ -112,13 +113,21 @@ def run_upgrade(self, from_kafka_version):
@matrix(from_kafka_version=[str(LATEST_3_1), str(LATEST_3_2), str(LATEST_3_3), str(LATEST_3_4), str(LATEST_3_5), str(DEV_BRANCH)],
use_new_coordinator=[True, False],
metadata_quorum=[combined_kraft])
def test_combined_mode_upgrade(self, from_kafka_version, metadata_quorum, use_new_coordinator=False):
self.run_upgrade(from_kafka_version)
@matrix(from_kafka_version=[str(LATEST_3_7), str(DEV_BRANCH)],
use_new_coordinator=[True],
metadata_quorum=[combined_kraft],
group_protocol=consumer_group.all_group_protocols)
def test_combined_mode_upgrade(self, from_kafka_version, metadata_quorum, use_new_coordinator=False, group_protocol=None):
self.run_upgrade(from_kafka_version, group_protocol)

@cluster(num_nodes=8)
@matrix(from_kafka_version=[str(LATEST_3_1), str(LATEST_3_2), str(LATEST_3_3), str(LATEST_3_4), str(LATEST_3_5), str(DEV_BRANCH)],
use_new_coordinator=[True, False],
@matrix(from_kafka_version=[str(LATEST_3_1), str(LATEST_3_2), str(LATEST_3_3), str(LATEST_3_4), str(LATEST_3_5), str(DEV_BRANCH)],
use_new_coordinator=[True, False],
metadata_quorum=[isolated_kraft])
def test_isolated_mode_upgrade(self, from_kafka_version, metadata_quorum, use_new_coordinator=False):
self.run_upgrade(from_kafka_version)
@matrix(from_kafka_version=[str(LATEST_3_7), str(DEV_BRANCH)],
use_new_coordinator=[True],
metadata_quorum=[isolated_kraft],
group_protocol=consumer_group.all_group_protocols)
def test_isolated_mode_upgrade(self, from_kafka_version, metadata_quorum, use_new_coordinator=False, group_protocol=None):
self.run_upgrade(from_kafka_version, group_protocol)