From 1804982a8017e26686c21cf7c06fd595a5176d28 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Mon, 1 May 2023 16:00:59 -0700 Subject: [PATCH 1/3] HOTFIX: fix broken Streams upgrade system test --- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- .../org/apache/kafka/streams/tests/StreamsUpgradeTest.java | 3 +-- 9 files changed, 9 insertions(+), 18 deletions(-) diff --git a/streams/upgrade-system-tests-24/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-24/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 9d08663d9b37f..06a363269031c 100644 --- a/streams/upgrade-system-tests-24/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-24/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -71,7 +70,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000L); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-25/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-25/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 69c46de37af14..efa32fc8be208 100644 --- a/streams/upgrade-system-tests-25/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-25/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -71,7 +70,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000L); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-26/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-26/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 0844552134a03..77a7cbbc3c036 100644 --- a/streams/upgrade-system-tests-26/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-26/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -71,7 +70,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000L); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-27/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-27/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 32d8d9408f57b..bf2e979815b89 100644 --- a/streams/upgrade-system-tests-27/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-27/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -71,7 +70,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000L); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-28/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-28/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index db17d73bcbaca..792359151327f 100644 --- a/streams/upgrade-system-tests-28/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-28/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -71,7 +70,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-30/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-30/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 0751516d76c9b..4951380d97958 100644 --- a/streams/upgrade-system-tests-30/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-30/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -73,7 +72,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-31/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-31/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 311d30ba40038..a622a33b8c971 100644 --- a/streams/upgrade-system-tests-31/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-31/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -19,7 +19,6 @@ import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; -import java.util.Random; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; @@ -73,7 +72,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-32/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-32/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 7419896a0bf2f..c5a279747049b 100644 --- a/streams/upgrade-system-tests-32/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-32/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -30,7 +30,6 @@ import org.apache.kafka.streams.processor.api.Record; import java.util.Properties; -import java.util.Random; import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; @@ -73,7 +72,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); config.putAll(streamsProperties); diff --git a/streams/upgrade-system-tests-33/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/upgrade-system-tests-33/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index 60b8305bc3587..4beab16fb3850 100644 --- a/streams/upgrade-system-tests-33/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/upgrade-system-tests-33/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -30,7 +30,6 @@ import org.apache.kafka.streams.processor.api.Record; import java.util.Properties; -import java.util.Random; import static org.apache.kafka.streams.tests.SmokeTestUtil.intSerde; import static org.apache.kafka.streams.tests.SmokeTestUtil.stringSerde; @@ -73,7 +72,7 @@ public static void main(final String[] args) throws Exception { final Properties config = new Properties(); config.setProperty( StreamsConfig.APPLICATION_ID_CONFIG, - "StreamsUpgradeTest-" + new Random().nextLong()); + "StreamsUpgradeTest"); config.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); config.putAll(streamsProperties); From a8f3172b9708794b34f48ac924fed2264d40fa07 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Tue, 2 May 2023 17:11:44 -0700 Subject: [PATCH 2/3] fix dev version --- tests/kafkatest/version.py | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/tests/kafkatest/version.py b/tests/kafkatest/version.py index c31adc256bab0..a285648ca5680 100644 --- a/tests/kafkatest/version.py +++ b/tests/kafkatest/version.py @@ -119,7 +119,7 @@ def get_version(node=None): return DEV_BRANCH DEV_BRANCH = KafkaVersion("dev") -DEV_VERSION = KafkaVersion("3.5.0-SNAPSHOT") +DEV_VERSION = KafkaVersion("3.6.0-SNAPSHOT") LATEST_METADATA_VERSION = "3.3" @@ -249,7 +249,3 @@ def get_version(node=None): # 3.5.x versions V_3_5_0 = KafkaVersion("3.5.0") LATEST_3_5 = V_3_5_0 - -# 3.6.x versions -V_3_6_0 = KafkaVersion("3.6.0") -LATEST_3_6 = V_3_6_0 From 4c082f4567f2bcf0b94307f28bae4473a66db02e Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Thu, 4 May 2023 16:10:49 -0700 Subject: [PATCH 3/3] fix test setup --- tests/kafkatest/tests/streams/streams_upgrade_test.py | 10 +++++++--- tests/kafkatest/version.py | 3 ++- 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/tests/kafkatest/tests/streams/streams_upgrade_test.py b/tests/kafkatest/tests/streams/streams_upgrade_test.py index cf276364ea457..9039d9d7473ab 100644 --- a/tests/kafkatest/tests/streams/streams_upgrade_test.py +++ b/tests/kafkatest/tests/streams/streams_upgrade_test.py @@ -38,9 +38,13 @@ str(DEV_BRANCH)] metadata_1_versions = [str(LATEST_0_10_0)] -metadata_2_versions = [str(LATEST_0_10_1), str(LATEST_0_10_2), str(LATEST_0_11_0), str(LATEST_1_0), str(LATEST_1_1)] -fk_join_versions = [str(LATEST_2_4), str(LATEST_2_5), str(LATEST_2_6), str(LATEST_2_7), str(LATEST_2_8), - str(LATEST_3_0), str(LATEST_3_1), str(LATEST_3_2), str(LATEST_3_3)] +metadata_2_versions = [str(LATEST_0_10_1), str(LATEST_0_10_2), str(LATEST_0_11_0), str(LATEST_1_0), str(LATEST_1_1), + str(LATEST_2_4), str(LATEST_2_5), str(LATEST_2_6), str(LATEST_2_7), str(LATEST_2_8), + str(LATEST_3_0)] +# upgrading from version (2.4...3.0) is broken and only fixed later in 3.1 +# we cannot test two bounce rolling upgrade because we know it's broken +# instead we add version 2.4...3.0 to the `metadata_2_versions` upgrade list +fk_join_versions = [str(LATEST_3_1), str(LATEST_3_2), str(LATEST_3_3)] """ After each release one should first check that the released version has been uploaded to diff --git a/tests/kafkatest/version.py b/tests/kafkatest/version.py index a285648ca5680..967674b551c5e 100644 --- a/tests/kafkatest/version.py +++ b/tests/kafkatest/version.py @@ -107,7 +107,8 @@ def supports_topic_ids_when_using_zk(self): return self >= V_2_8_0 def supports_fk_joins(self): - return hasattr(self, "version") and self >= V_2_4_0 + # while we support FK joins since 2.4, rolling upgrade is broken in older versions and only fixed in 3.1 + return hasattr(self, "version") and self >= V_3_1_2 def get_version(node=None): """Return the version attached to the given node.