From 0618d89546c744c05df19b6b6077bf38575e69f5 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Tue, 20 Dec 2022 18:30:02 +0200 Subject: [PATCH 1/4] [improve][io] Upgrade Kafka client, connect runtime to 2.8.2 and Confluent version to 6.2.8 - Confluent 6.2.x is compatible with Kafka client 2.8.2 --- pom.xml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/pom.xml b/pom.xml index 5c4b55cdc77a6..49dbc22a1d2a5 100644 --- a/pom.xml +++ b/pom.xml @@ -160,7 +160,7 @@ flexible messaging model and an intuitive client API. 2.2.0 3.11.2 4.4.20 - 2.7.2 + 2.8.2 5.5.3 1.12.262 1.10.2 @@ -187,9 +187,9 @@ flexible messaging model and an intuitive client API. 31.0.1-jre 1.0 0.16.1 - 7.0.1 - 5.3.0 - 5.3.0 + 6.2.8 + ${confluent.version} + ${confluent.version} 0.20 2.12.1 1.82 From 06af7486fe10eb117eec7fb7b37d6a4fb1d2dbe3 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Tue, 20 Dec 2022 22:54:01 +0200 Subject: [PATCH 2/4] Replace kafka.confluent.schemaregistryclient.version and kafka.confluent.avroserializer.version with confluent.version --- pom.xml | 2 -- pulsar-io/kafka/pom.xml | 4 ++-- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/pom.xml b/pom.xml index 49dbc22a1d2a5..77ae1053790ca 100644 --- a/pom.xml +++ b/pom.xml @@ -188,8 +188,6 @@ flexible messaging model and an intuitive client API. 1.0 0.16.1 6.2.8 - ${confluent.version} - ${confluent.version} 0.20 2.12.1 1.82 diff --git a/pulsar-io/kafka/pom.xml b/pulsar-io/kafka/pom.xml index 2e6e24d3c5d46..ec8de4c2e2e12 100644 --- a/pulsar-io/kafka/pom.xml +++ b/pulsar-io/kafka/pom.xml @@ -84,13 +84,13 @@ io.confluent kafka-schema-registry-client - ${kafka.confluent.schemaregistryclient.version} + ${confluent.version} io.confluent kafka-avro-serializer - ${kafka.confluent.avroserializer.version} + ${confluent.version} From b148a1208574fc797d2b06e06d5d17e8bbd6f297 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Tue, 20 Dec 2022 22:54:30 +0200 Subject: [PATCH 3/4] Use confluent version in AvroKafkaSourceTest --- tests/integration/pom.xml | 2 +- .../tests/integration/io/sources/AvroKafkaSourceTest.java | 5 +++-- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/tests/integration/pom.xml b/tests/integration/pom.xml index a9158f0b395c4..652fc7543500e 100644 --- a/tests/integration/pom.xml +++ b/tests/integration/pom.xml @@ -280,7 +280,7 @@ maven-surefire-plugin ${testJacocoAgentArgument} -XX:+ExitOnOutOfMemoryError -Xmx1G -XX:MaxDirectMemorySize=1G - -Dio.netty.leakDetectionLevel=advanced + -Dio.netty.leakDetectionLevel=advanced -Dconfluent.version=${confluent.version} ${test.additional.args} false diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/AvroKafkaSourceTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/AvroKafkaSourceTest.java index 5fcf4bd284d28..913b4e376748f 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/AvroKafkaSourceTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/AvroKafkaSourceTest.java @@ -71,6 +71,7 @@ */ @Slf4j public class AvroKafkaSourceTest extends PulsarFunctionsTestBase { + public static final String CONFLUENT_PLATFORM_VERSION = System.getProperty("confluent.version", "6.2.8"); private static final String SOURCE_TYPE = "kafka"; @@ -139,7 +140,8 @@ public String getBootstrapServers() { } protected EnhancedKafkaContainer createKafkaContainer(PulsarCluster cluster) { - return (EnhancedKafkaContainer) new EnhancedKafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:6.0.1")) + return (EnhancedKafkaContainer) new EnhancedKafkaContainer( + DockerImageName.parse("confluentinc/cp-kafka:" + CONFLUENT_PLATFORM_VERSION)) .withEmbeddedZookeeper() .withCreateContainerCmdModifier(createContainerCmd -> createContainerCmd .withName(kafkaContainerName) @@ -480,7 +482,6 @@ protected void getSourceInfoNotFound(String tenant, String namespace, String sou } public class SchemaRegistryContainer extends GenericContainer { - public static final String CONFLUENT_PLATFORM_VERSION = "6.0.1"; private static final int SCHEMA_REGISTRY_INTERNAL_PORT = 8081; public SchemaRegistryContainer(String boostrapServers) throws Exception { From 2cca77d5d533b9872f10329fc787fcc6c1f666f6 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Tue, 20 Dec 2022 22:59:58 +0200 Subject: [PATCH 4/4] Use confluent.version in KafkaSinkTester --- .../pulsar/tests/integration/io/sinks/KafkaSinkTester.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java index ceb0c820b629f..d474efaadfeb7 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/KafkaSinkTester.java @@ -35,12 +35,14 @@ import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.testcontainers.containers.Container.ExecResult; import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.utility.DockerImageName; /** * A tester for testing kafka sink. */ @Slf4j public class KafkaSinkTester extends SinkTester { + public static final String CONFLUENT_PLATFORM_VERSION = System.getProperty("confluent.version", "6.2.8"); private final String kafkaTopicName; private KafkaConsumer kafkaConsumer; @@ -63,7 +65,7 @@ public KafkaSinkTester(String containerName) { @SuppressWarnings("deprecation") @Override protected KafkaContainer createSinkService(PulsarCluster cluster) { - return new KafkaContainer() + return new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:" + CONFLUENT_PLATFORM_VERSION)) .withEmbeddedZookeeper() .withNetworkAliases(containerName) .withCreateContainerCmdModifier(createContainerCmd -> createContainerCmd