This is an automated email from the ASF dual-hosted git repository.

Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new ec397ce5003 [KafkaIO] Remove build support for Kafka clients before 
3.9.2 (#39284)
ec397ce5003 is described below

commit ec397ce5003b40f05aaa1a56683211aaf942905c
Author: Steven van Rossum <[email protected]>
AuthorDate: Mon Jul 20 20:32:50 2026 +0200

    [KafkaIO] Remove build support for Kafka clients before 3.9.2 (#39284)
    
    * Remove support for Kafka clients older than 3.9.2
    
    * Resolve capability conflicts for Flink 1.x, Flink 2.0, Spark 3 and Spark 4
    
    * Resolve capability conflicts for load tests and watermarks
    
    * Replace obsolete signature of overridden method close in mock consumer
    
    * Replace obsolete class KafkaServerStartable with KafkaServer
    
    * Fix type ambiguity of method argument in test
    
    * Handle Kafka and executor timeouts the same
---
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |  6 +++--
 examples/java/build.gradle                         |  4 ++--
 examples/java/common.gradle                        |  1 +
 .../beam/it/kafka/KafkaResourceManagerTest.java    |  3 ++-
 runners/flink/1.19/build.gradle                    |  6 +++++
 runners/flink/1.19/job-server/build.gradle         |  6 +++++
 runners/flink/1.20/build.gradle                    |  6 +++++
 runners/flink/1.20/job-server/build.gradle         |  6 +++++
 runners/flink/2.0/build.gradle                     |  6 +++++
 runners/flink/2.0/job-server/build.gradle          |  6 +++++
 runners/flink/2.1/build.gradle                     |  7 ------
 runners/flink/2.1/job-server/build.gradle          |  7 ------
 runners/flink/2.2/build.gradle                     |  7 ------
 runners/flink/2.2/job-server/build.gradle          |  7 ------
 runners/spark/3/build.gradle                       |  6 +++++
 runners/spark/3/job-server/build.gradle            |  8 ++++++-
 runners/spark/4/build.gradle                       |  5 ++++
 runners/spark/4/job-server/build.gradle            |  6 +++++
 runners/spark/spark_runner.gradle                  |  5 ++--
 .../streaming/utils/EmbeddedKafkaCluster.java      | 17 ++++++++-----
 sdks/java/io/kafka/build.gradle                    | 12 +++-------
 sdks/java/io/kafka/kafka-201/build.gradle          | 24 -------------------
 sdks/java/io/kafka/kafka-231/build.gradle          | 24 -------------------
 sdks/java/io/kafka/kafka-241/build.gradle          | 24 -------------------
 sdks/java/io/kafka/kafka-282/build.gradle          | 24 -------------------
 sdks/java/io/kafka/kafka-312/build.gradle          | 24 -------------------
 sdks/java/io/kafka/kafka-390/build.gradle          | 24 -------------------
 .../io/kafka/{kafka-251 => kafka-392}/build.gradle |  6 ++---
 .../beam/sdk/io/kafka/KafkaUnboundedReader.java    | 28 +++++++++++++---------
 .../beam/sdk/io/kafka/KafkaCommitOffsetTest.java   |  4 ++--
 sdks/java/testing/kafka-service/build.gradle       |  3 ++-
 .../apache/beam/sdk/testing/kafka/LocalKafka.java  | 11 ++++++---
 sdks/java/testing/load-tests/build.gradle          | 16 +++++++++++++
 sdks/java/testing/nexmark/build.gradle             |  7 +++---
 sdks/java/testing/tpcds/build.gradle               |  7 +++---
 sdks/java/testing/watermarks/build.gradle          | 25 ++++++++++++++++++-
 settings.gradle.kts                                | 16 ++-----------
 37 files changed, 167 insertions(+), 237 deletions(-)

diff --git 
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy 
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index 3e7ffa89b74..d73f2e7a2ba 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -639,7 +639,7 @@ class BeamModulePlugin implements Plugin<Project> {
     def jaxb_api_version = "2.3.3"
     def jsr305_version = "3.0.2"
     def everit_json_version = "1.14.2"
-    def kafka_version = "2.4.1"
+    def kafka_version = "3.9.2"
     def log4j2_version = "2.25.4"
     def nemo_version = "0.1"
     // [bomupgrader] determined by: io.grpc:grpc-netty, consistent with: 
google_cloud_platform_libraries_bom
@@ -850,8 +850,10 @@ class BeamModulePlugin implements Plugin<Project> {
         jupiter_api                                 : 
"org.junit.jupiter:junit-jupiter-api:$jupiter_version",
         jupiter_engine                              : 
"org.junit.jupiter:junit-jupiter-engine:$jupiter_version",
         jupiter_params                              : 
"org.junit.jupiter:junit-jupiter-params:$jupiter_version",
-        kafka                                       : 
"org.apache.kafka:kafka_2.11:$kafka_version",
+        kafka_scala_2_12                            : 
"org.apache.kafka:kafka_2.12:$kafka_version",
+        kafka_scala_2_13                            : 
"org.apache.kafka:kafka_2.13:$kafka_version",
         kafka_clients                               : 
"org.apache.kafka:kafka-clients:$kafka_version",
+        kafka_server                                : 
"org.apache.kafka:kafka-server:$kafka_version",
         log4j                                       : "log4j:log4j:1.2.17",
         log4j_over_slf4j                            : 
"org.slf4j:log4j-over-slf4j:$slf4j_version",
         log4j2_api                                  : 
"org.apache.logging.log4j:log4j-api:$log4j2_version",
diff --git a/examples/java/build.gradle b/examples/java/build.gradle
index 34d884b9778..84ace728362 100644
--- a/examples/java/build.gradle
+++ b/examples/java/build.gradle
@@ -44,8 +44,8 @@ dependencies {
   if (project.findProperty('testJavaVersion') == '21' || 
JavaVersion.current().compareTo(JavaVersion.VERSION_21) >= 0) {
     // this dependency is a provided dependency for kafka-avro-serializer. It 
is not needed to compile with Java<=17
     // but needed for compile only under Java21, specifically, required for 
extending from AbstractKafkaAvroDeserializer
-    compileOnly library.java.kafka
-    permitUnusedDeclared library.java.kafka
+    compileOnly library.java.kafka_scala_2_12
+    permitUnusedDeclared library.java.kafka_scala_2_12
   }
   implementation library.java.kafka_clients
   implementation project(path: ":sdks:java:core", configuration: "shadow")
diff --git a/examples/java/common.gradle b/examples/java/common.gradle
index f32667d733f..91d07aff76b 100644
--- a/examples/java/common.gradle
+++ b/examples/java/common.gradle
@@ -36,6 +36,7 @@ configurations.sparkRunnerPreCommit {
   exclude group: "org.slf4j", module: "slf4j-jdk14"
 }
 resolveCapabilitiesConflict(configurations.flinkRunnerPreCommit, 
'org.lz4:lz4-java', 'at.yawk.lz4')
+resolveCapabilitiesConflict(configurations.sparkRunnerPreCommit, 
'org.lz4:lz4-java', 'at.yawk.lz4')
 
 
 dependencies {
diff --git 
a/it/kafka/src/test/java/org/apache/beam/it/kafka/KafkaResourceManagerTest.java 
b/it/kafka/src/test/java/org/apache/beam/it/kafka/KafkaResourceManagerTest.java
index 8c870815efc..f25bca06c22 100644
--- 
a/it/kafka/src/test/java/org/apache/beam/it/kafka/KafkaResourceManagerTest.java
+++ 
b/it/kafka/src/test/java/org/apache/beam/it/kafka/KafkaResourceManagerTest.java
@@ -166,7 +166,8 @@ public final class KafkaResourceManagerTest {
     KafkaResourceManager tm = new KafkaResourceManager(kafkaClient, container, 
builder);
 
     tm.cleanupAll();
-    verify(kafkaClient).deleteTopics(argThat(list -> list.size() == 
numTopics));
+    verify(kafkaClient)
+        .deleteTopics(argThat((Collection<String> list) -> list.size() == 
numTopics));
   }
 
   @Test
diff --git a/runners/flink/1.19/build.gradle b/runners/flink/1.19/build.gradle
index 1545da25847..4116b817eec 100644
--- a/runners/flink/1.19/build.gradle
+++ b/runners/flink/1.19/build.gradle
@@ -23,3 +23,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "../flink_runner.gradle"
+
+// Flink 1.19 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/flink/1.19/job-server/build.gradle 
b/runners/flink/1.19/job-server/build.gradle
index 332f04e08ce..c9e09a5de8c 100644
--- a/runners/flink/1.19/job-server/build.gradle
+++ b/runners/flink/1.19/job-server/build.gradle
@@ -29,3 +29,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "$basePath/flink_job_server.gradle"
+
+// Flink 1.19 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/flink/1.20/build.gradle b/runners/flink/1.20/build.gradle
index 4c148321ed4..86a1cb7ef54 100644
--- a/runners/flink/1.20/build.gradle
+++ b/runners/flink/1.20/build.gradle
@@ -23,3 +23,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "../flink_runner.gradle"
+
+// Flink 1.20 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/flink/1.20/job-server/build.gradle 
b/runners/flink/1.20/job-server/build.gradle
index e5fdd1febf9..9f129b2c021 100644
--- a/runners/flink/1.20/job-server/build.gradle
+++ b/runners/flink/1.20/job-server/build.gradle
@@ -29,3 +29,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "$basePath/flink_job_server.gradle"
+
+// Flink 1.20 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/flink/2.0/build.gradle b/runners/flink/2.0/build.gradle
index 490bc593f40..4e034173fb9 100644
--- a/runners/flink/2.0/build.gradle
+++ b/runners/flink/2.0/build.gradle
@@ -41,3 +41,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "../flink_runner.gradle"
+
+// Flink 2.0 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/flink/2.0/job-server/build.gradle 
b/runners/flink/2.0/job-server/build.gradle
index 6d068f83949..ebac0a8e039 100644
--- a/runners/flink/2.0/job-server/build.gradle
+++ b/runners/flink/2.0/job-server/build.gradle
@@ -29,3 +29,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "$basePath/flink_job_server.gradle"
+
+// Flink 2.0 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/flink/2.1/build.gradle b/runners/flink/2.1/build.gradle
index 1e0d565b50d..32c71eb9648 100644
--- a/runners/flink/2.1/build.gradle
+++ b/runners/flink/2.1/build.gradle
@@ -41,10 +41,3 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "../flink_runner.gradle"
-
-// Flink 2.1 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
-// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
-configurations.all {
-  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
-}
-
diff --git a/runners/flink/2.1/job-server/build.gradle 
b/runners/flink/2.1/job-server/build.gradle
index 0910fef1120..3135abcb19b 100644
--- a/runners/flink/2.1/job-server/build.gradle
+++ b/runners/flink/2.1/job-server/build.gradle
@@ -29,10 +29,3 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "$basePath/flink_job_server.gradle"
-
-// Flink 2.1 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
-// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
-configurations.all {
-  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
-}
-
diff --git a/runners/flink/2.2/build.gradle b/runners/flink/2.2/build.gradle
index 0321dcf42d1..1e1d7fdd5e7 100644
--- a/runners/flink/2.2/build.gradle
+++ b/runners/flink/2.2/build.gradle
@@ -56,10 +56,3 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "../flink_runner.gradle"
-
-// Flink 2.2 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
-// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
-configurations.all {
-  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
-}
-
diff --git a/runners/flink/2.2/job-server/build.gradle 
b/runners/flink/2.2/job-server/build.gradle
index f116f5a1dcb..52724532921 100644
--- a/runners/flink/2.2/job-server/build.gradle
+++ b/runners/flink/2.2/job-server/build.gradle
@@ -29,10 +29,3 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "$basePath/flink_job_server.gradle"
-
-// Flink 2.2 uses at.yawk.lz4:lz4-java instead of org.lz4:lz4-java
-// Explicitly prefer Flink's at.yawk.lz4 version to resolve capability conflict
-configurations.all {
-  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
-}
-
diff --git a/runners/spark/3/build.gradle b/runners/spark/3/build.gradle
index c2c492a6b2e..274fa45d2f7 100644
--- a/runners/spark/3/build.gradle
+++ b/runners/spark/3/build.gradle
@@ -88,3 +88,9 @@ tasks.register("sparkVersionsTest") {
   group = "Verification"
   dependsOn sparkVersions.collect{k,v -> "sparkVersion${k}Test"}
 }
+
+// Spark 3 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/spark/3/job-server/build.gradle 
b/runners/spark/3/job-server/build.gradle
index 68bb8d9a10e..59165e80be5 100644
--- a/runners/spark/3/job-server/build.gradle
+++ b/runners/spark/3/job-server/build.gradle
@@ -28,4 +28,10 @@ project.ext {
 }
 
 // Load the main build script which contains all build logic.
-apply from: "$basePath/spark_job_server.gradle"
\ No newline at end of file
+apply from: "$basePath/spark_job_server.gradle"
+
+// Spark 3 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/spark/4/build.gradle b/runners/spark/4/build.gradle
index ec1af8df38a..f2746b06158 100644
--- a/runners/spark/4/build.gradle
+++ b/runners/spark/4/build.gradle
@@ -61,3 +61,8 @@ tasks.named("copyTestSourceOverrides") {
   exclude "**/translation/streaming/**"
 }
 
+// Spark 4 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/spark/4/job-server/build.gradle 
b/runners/spark/4/job-server/build.gradle
index 598cf3b4913..3154d92ea32 100644
--- a/runners/spark/4/job-server/build.gradle
+++ b/runners/spark/4/job-server/build.gradle
@@ -29,3 +29,9 @@ project.ext {
 
 // Load the main build script which contains all build logic.
 apply from: "$basePath/spark_job_server.gradle"
+
+// Spark 4 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
+// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
+configurations.all {
+  resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
+}
diff --git a/runners/spark/spark_runner.gradle 
b/runners/spark/spark_runner.gradle
index 1da044ab7b5..77da3d36db9 100644
--- a/runners/spark/spark_runner.gradle
+++ b/runners/spark/spark_runner.gradle
@@ -287,10 +287,9 @@ dependencies {
   testImplementation project(path: ":sdks:java:extensions:avro", 
configuration: "testRuntimeMigration")
   testImplementation project(":sdks:java:harness")
   testImplementation library.java.avro
-  // kafka_2.13 artifacts were first published in 2.5.0; use a later version 
for Scala 2.13
-  def kafka_version = (spark_scala_version == '2.13') ? '2.8.0' : '2.4.1'
-  testImplementation 
"org.apache.kafka:kafka_$spark_scala_version:$kafka_version"
+  testImplementation (spark_scala_version == '2.13' ? 
library.java.kafka_scala_2_13 : library.java.kafka_scala_2_12)
   testImplementation library.java.kafka_clients
+  testImplementation library.java.kafka_server
   testImplementation library.java.junit
   testImplementation library.java.mockito_core
   testImplementation "org.assertj:assertj-core:3.11.1"
diff --git 
a/runners/spark/src/test/java/org/apache/beam/runners/spark/translation/streaming/utils/EmbeddedKafkaCluster.java
 
b/runners/spark/src/test/java/org/apache/beam/runners/spark/translation/streaming/utils/EmbeddedKafkaCluster.java
index df5646fed59..3acbf664e2c 100644
--- 
a/runners/spark/src/test/java/org/apache/beam/runners/spark/translation/streaming/utils/EmbeddedKafkaCluster.java
+++ 
b/runners/spark/src/test/java/org/apache/beam/runners/spark/translation/streaming/utils/EmbeddedKafkaCluster.java
@@ -29,7 +29,7 @@ import java.util.List;
 import java.util.Properties;
 import java.util.Random;
 import kafka.server.KafkaConfig;
-import kafka.server.KafkaServerStartable;
+import kafka.server.KafkaServer;
 import org.apache.zookeeper.server.NIOServerCnxnFactory;
 import org.apache.zookeeper.server.ServerCnxnFactory;
 import org.apache.zookeeper.server.ZooKeeperServer;
@@ -47,7 +47,7 @@ public class EmbeddedKafkaCluster {
 
   private final String brokerList;
 
-  private final List<KafkaServerStartable> brokers;
+  private final List<KafkaServer> brokers;
   private final List<File> logDirs;
 
   private EmbeddedKafkaCluster(String zkConnection) {
@@ -114,15 +114,20 @@ public class EmbeddedKafkaCluster {
       properties.setProperty("offsets.topic.replication.factor", "1");
       properties.setProperty("log.flush.interval.messages", String.valueOf(1));
 
-      KafkaServerStartable broker = startBroker(properties);
+      KafkaServer broker = startBroker(properties);
 
       brokers.add(broker);
       logDirs.add(logDir);
     }
   }
 
-  private static KafkaServerStartable startBroker(Properties props) {
-    KafkaServerStartable server = new KafkaServerStartable(new 
KafkaConfig(props));
+  private static KafkaServer startBroker(Properties props) {
+    KafkaServer server =
+        new KafkaServer(
+            new KafkaConfig(props),
+            KafkaServer.$lessinit$greater$default$2(),
+            KafkaServer.$lessinit$greater$default$3(),
+            KafkaServer.$lessinit$greater$default$4());
     server.startup();
     return server;
   }
@@ -148,7 +153,7 @@ public class EmbeddedKafkaCluster {
 
   @SuppressWarnings("Slf4jDoNotLogMessageOfExceptionExplicitly")
   public void shutdown() {
-    for (KafkaServerStartable broker : brokers) {
+    for (KafkaServer broker : brokers) {
       try {
         broker.shutdown();
       } catch (Exception e) {
diff --git a/sdks/java/io/kafka/build.gradle b/sdks/java/io/kafka/build.gradle
index 13969f7a9ae..0d28469eae5 100644
--- a/sdks/java/io/kafka/build.gradle
+++ b/sdks/java/io/kafka/build.gradle
@@ -36,13 +36,7 @@ ext {
 }
 
 def kafkaVersions = [
-    '201': "2.0.1",
-    '231': "2.3.1",
-    '241': "2.4.1",
-    '251': "2.5.1",
-    '282': "2.8.2",
-    '312': "3.1.2",
-    '390': "3.9.0",
+    '392': "3.9.2",
 ]
 
 kafkaVersions.each{k,v -> configurations.create("kafkaVersion$k")}
@@ -63,8 +57,8 @@ dependencies {
   if (JavaVersion.current().compareTo(JavaVersion.VERSION_21) >= 0) {
     // this dependency is a provided dependency for kafka-avro-serializer. It 
is not needed to compile with Java<=17
     // but needed for compile only under Java21, specifically, required for 
extending from AbstractKafkaAvroDeserializer
-    compileOnly library.java.kafka
-    permitUnusedDeclared library.java.kafka
+    compileOnly library.java.kafka_scala_2_12
+    permitUnusedDeclared library.java.kafka_scala_2_12
   }
   testImplementation library.java.kafka_clients
   testImplementation project(path: ":runners:core-java")
diff --git a/sdks/java/io/kafka/kafka-201/build.gradle 
b/sdks/java/io/kafka/kafka-201/build.gradle
deleted file mode 100644
index a26ca4ac19c..00000000000
--- a/sdks/java/io/kafka/kafka-201/build.gradle
+++ /dev/null
@@ -1,24 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * License); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an AS IS BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-project.ext {
-    delimited="2.0.1"
-    undelimited="201"
-    sdfCompatible=true
-}
-
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
diff --git a/sdks/java/io/kafka/kafka-231/build.gradle 
b/sdks/java/io/kafka/kafka-231/build.gradle
deleted file mode 100644
index 712158dcd3a..00000000000
--- a/sdks/java/io/kafka/kafka-231/build.gradle
+++ /dev/null
@@ -1,24 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * License); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an AS IS BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-project.ext {
-    delimited="2.3.1"
-    undelimited="231"
-    sdfCompatible=true
-}
-
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
diff --git a/sdks/java/io/kafka/kafka-241/build.gradle 
b/sdks/java/io/kafka/kafka-241/build.gradle
deleted file mode 100644
index c0ac7df674b..00000000000
--- a/sdks/java/io/kafka/kafka-241/build.gradle
+++ /dev/null
@@ -1,24 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * License); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an AS IS BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-project.ext {
-    delimited="2.4.1"
-    undelimited="241"
-    sdfCompatible=true
-}
-
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
diff --git a/sdks/java/io/kafka/kafka-282/build.gradle 
b/sdks/java/io/kafka/kafka-282/build.gradle
deleted file mode 100644
index b754d93077e..00000000000
--- a/sdks/java/io/kafka/kafka-282/build.gradle
+++ /dev/null
@@ -1,24 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * License); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an AS IS BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-project.ext {
-    delimited="2.8.2"
-    undelimited="282"
-    sdfCompatible=true
-}
-
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
diff --git a/sdks/java/io/kafka/kafka-312/build.gradle 
b/sdks/java/io/kafka/kafka-312/build.gradle
deleted file mode 100644
index af2ad3717b6..00000000000
--- a/sdks/java/io/kafka/kafka-312/build.gradle
+++ /dev/null
@@ -1,24 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * License); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an AS IS BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-project.ext {
-    delimited="3.1.2"
-    undelimited="312"
-    sdfCompatible=true
-}
-
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
diff --git a/sdks/java/io/kafka/kafka-390/build.gradle 
b/sdks/java/io/kafka/kafka-390/build.gradle
deleted file mode 100644
index 8c882138626..00000000000
--- a/sdks/java/io/kafka/kafka-390/build.gradle
+++ /dev/null
@@ -1,24 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * License); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an AS IS BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-project.ext {
-    delimited="3.9.0"
-    undelimited="390"
-    sdfCompatible=true
-}
-
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
diff --git a/sdks/java/io/kafka/kafka-251/build.gradle 
b/sdks/java/io/kafka/kafka-392/build.gradle
similarity index 90%
rename from sdks/java/io/kafka/kafka-251/build.gradle
rename to sdks/java/io/kafka/kafka-392/build.gradle
index 4de9f97a738..3df1ebeb69a 100644
--- a/sdks/java/io/kafka/kafka-251/build.gradle
+++ b/sdks/java/io/kafka/kafka-392/build.gradle
@@ -16,9 +16,9 @@
  * limitations under the License.
  */
 project.ext {
-    delimited="2.5.1"
-    undelimited="251"
+    delimited="3.9.2"
+    undelimited="392"
     sdfCompatible=true
 }
 
-apply from: "../kafka-integration-test.gradle"
\ No newline at end of file
+apply from: "../kafka-integration-test.gradle"
diff --git 
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaUnboundedReader.java
 
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaUnboundedReader.java
index c5dc5c408fe..c3aa6f5418b 100644
--- 
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaUnboundedReader.java
+++ 
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaUnboundedReader.java
@@ -30,6 +30,7 @@ import java.util.Map;
 import java.util.NoSuchElementException;
 import java.util.Optional;
 import java.util.Set;
+import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
@@ -114,17 +115,22 @@ class KafkaUnboundedReader<K, V> extends 
UnboundedReader<KafkaRecord<K, V>> {
       try {
         Duration timeout = resolveDefaultApiTimeout(spec);
         future.get(timeout.getMillis(), TimeUnit.MILLISECONDS);
-      } catch (TimeoutException e) {
-        consumer.wakeup(); // This unblocks consumer stuck on network I/O.
-        // Likely reason : Kafka servers are configured to advertise internal 
ips, but
-        // those ips are not accessible from workers outside.
-        String msg =
-            String.format(
-                "%s: Timeout while initializing partition '%s'. "
-                    + "Kafka client may not be able to connect to servers.",
-                this, pState.topicPartition);
-        LOG.error("{}", msg);
-        throw new IOException(msg);
+      } catch (TimeoutException | ExecutionException e) {
+        if (e instanceof TimeoutException
+            || e.getCause() instanceof 
org.apache.kafka.common.errors.TimeoutException) {
+          // TODO: Find out if manually waking up was only relevant for legacy 
Kafka clients.
+          consumer.wakeup(); // This unblocks consumer stuck on network I/O.
+          // Likely reason : Kafka servers are configured to advertise 
internal ips, but
+          // those ips are not accessible from workers outside.
+          String msg =
+              String.format(
+                  "%s: Timeout while initializing partition '%s'. "
+                      + "Kafka client may not be able to connect to servers.",
+                  this, pState.topicPartition);
+          LOG.error("{}", msg);
+          throw new IOException(msg);
+        }
+        throw new IOException(e);
       } catch (Exception e) {
         throw new IOException(e);
       }
diff --git 
a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffsetTest.java
 
b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffsetTest.java
index c16e25510ab..b64f5edabe7 100644
--- 
a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffsetTest.java
+++ 
b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaCommitOffsetTest.java
@@ -17,11 +17,11 @@
  */
 package org.apache.beam.sdk.io.kafka;
 
+import java.time.Duration;
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
-import java.util.concurrent.TimeUnit;
 import org.apache.beam.sdk.coders.CannotProvideCoderException;
 import org.apache.beam.sdk.coders.KvCoder;
 import org.apache.beam.sdk.coders.StringUtf8Coder;
@@ -236,7 +236,7 @@ public class KafkaCommitOffsetTest {
     }
 
     @Override
-    public synchronized void close(long timeout, TimeUnit unit) {
+    public synchronized void close(Duration timeout) {
       // Ignore closing since we're using a single consumer.
     }
   }
diff --git a/sdks/java/testing/kafka-service/build.gradle 
b/sdks/java/testing/kafka-service/build.gradle
index abd186f98b1..1148ebe6572 100644
--- a/sdks/java/testing/kafka-service/build.gradle
+++ b/sdks/java/testing/kafka-service/build.gradle
@@ -27,7 +27,8 @@ ext.summary = """Self-contained Kafka service for testing IO 
transforms."""
 
 
 dependencies {
-  testImplementation library.java.kafka
+  testImplementation library.java.kafka_scala_2_12
+  testImplementation library.java.kafka_server
   testImplementation "org.apache.zookeeper:zookeeper:3.5.6"
   testRuntimeOnly library.java.slf4j_log4j12
 }
diff --git 
a/sdks/java/testing/kafka-service/src/test/java/org/apache/beam/sdk/testing/kafka/LocalKafka.java
 
b/sdks/java/testing/kafka-service/src/test/java/org/apache/beam/sdk/testing/kafka/LocalKafka.java
index 71ec61a3e41..dc88f3fe0a5 100644
--- 
a/sdks/java/testing/kafka-service/src/test/java/org/apache/beam/sdk/testing/kafka/LocalKafka.java
+++ 
b/sdks/java/testing/kafka-service/src/test/java/org/apache/beam/sdk/testing/kafka/LocalKafka.java
@@ -20,10 +20,10 @@ package org.apache.beam.sdk.testing.kafka;
 import java.nio.file.Files;
 import java.util.Properties;
 import kafka.server.KafkaConfig;
-import kafka.server.KafkaServerStartable;
+import kafka.server.KafkaServer;
 
 public class LocalKafka {
-  private final KafkaServerStartable server;
+  private final KafkaServer server;
 
   LocalKafka(int kafkaPort, int zookeeperPort) throws Exception {
     Properties kafkaProperties = new Properties();
@@ -31,7 +31,12 @@ public class LocalKafka {
     kafkaProperties.setProperty("zookeeper.connect", 
String.format("localhost:%s", zookeeperPort));
     kafkaProperties.setProperty("offsets.topic.replication.factor", "1");
     kafkaProperties.setProperty("log.dir", 
Files.createTempDirectory("kafka-log-").toString());
-    server = new KafkaServerStartable(KafkaConfig.fromProps(kafkaProperties));
+    server =
+        new KafkaServer(
+            KafkaConfig.fromProps(kafkaProperties),
+            KafkaServer.$lessinit$greater$default$2(),
+            KafkaServer.$lessinit$greater$default$3(),
+            KafkaServer.$lessinit$greater$default$4());
   }
 
   public void start() {
diff --git a/sdks/java/testing/load-tests/build.gradle 
b/sdks/java/testing/load-tests/build.gradle
index 699261963e6..ac3e391866d 100644
--- a/sdks/java/testing/load-tests/build.gradle
+++ b/sdks/java/testing/load-tests/build.gradle
@@ -40,6 +40,17 @@ def runnerDependency = (project.hasProperty(runnerProperty)
 def loadTestRunnerVersionProperty = "runner.version"
 def loadTestRunnerVersion = project.findProperty(loadTestRunnerVersionProperty)
 def isSparkRunner = runnerDependency.startsWith(":runners:spark:")
+def gtFlink20 = runnerDependency.startsWith(":runners:flink:") && {
+  def version = runnerDependency.substring(":runners:flink:".length())
+  try {
+    def parts = version.split('\\.')
+    def major = parts[0].toInteger()
+    def minor = parts.length > 1 ? parts[1].toInteger() : 0
+    return (major == 2 && minor > 0)
+  } catch (Exception e) {
+    return false
+  }
+}()
 def isDataflowRunner = 
":runners:google-cloud-dataflow-java".equals(runnerDependency)
 def isDataflowRunnerV2 = isDataflowRunner && "V2".equals(loadTestRunnerVersion)
 def runnerConfiguration = ":runners:direct-java".equals(runnerDependency) ? 
"shadow" : null
@@ -89,11 +100,16 @@ dependencies {
   gradleRun project(path: runnerDependency, configuration: runnerConfiguration)
 }
 
+if (!gtFlink20) {
+  resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
+}
+
 if (isSparkRunner) {
   configurations.gradleRun {
     // Using Spark runner causes a StackOverflowError if slf4j-jdk14 is on the 
classpath
     exclude group: "org.slf4j", module: "slf4j-jdk14"
   }
+  resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
 }
 
 def sparkJvmArgs() {
diff --git a/sdks/java/testing/nexmark/build.gradle 
b/sdks/java/testing/nexmark/build.gradle
index b554e9d9297..0eeaf931a88 100644
--- a/sdks/java/testing/nexmark/build.gradle
+++ b/sdks/java/testing/nexmark/build.gradle
@@ -39,13 +39,13 @@ def nexmarkRunnerDependency = 
project.findProperty(nexmarkRunnerProperty)
 def nexmarkRunnerVersionProperty = "nexmark.runner.version"
 def nexmarkRunnerVersion = project.findProperty(nexmarkRunnerVersionProperty)
 def isSparkRunner = nexmarkRunnerDependency.startsWith(":runners:spark:")
-def isFlink2 = nexmarkRunnerDependency.startsWith(":runners:flink:") && {
+def gtFlink20 = nexmarkRunnerDependency.startsWith(":runners:flink:") && {
   def version = nexmarkRunnerDependency.substring(":runners:flink:".length())
   try {
     def parts = version.split('\\.')
     def major = parts[0].toInteger()
     def minor = parts.length > 1 ? parts[1].toInteger() : 0
-    return (major > 2) || (major == 2 && minor >= 1)
+    return (major == 2 && minor > 0)
   } catch (Exception e) {
     return false
   }
@@ -103,7 +103,7 @@ dependencies {
   gradleRun project(path: nexmarkRunnerDependency, configuration: 
runnerConfiguration)
 }
 
-if (isFlink2) {
+if (!gtFlink20) {
   resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
 }
 
@@ -112,6 +112,7 @@ if (isSparkRunner) {
     // Using Spark runner causes a StackOverflowError if slf4j-jdk14 is on the 
classpath
     exclude group: "org.slf4j", module: "slf4j-jdk14"
   }
+  resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
 }
 
 def sparkJvmArgs() {
diff --git a/sdks/java/testing/tpcds/build.gradle 
b/sdks/java/testing/tpcds/build.gradle
index 15fdd480f07..60c2f8bfdd8 100644
--- a/sdks/java/testing/tpcds/build.gradle
+++ b/sdks/java/testing/tpcds/build.gradle
@@ -34,13 +34,13 @@ def tpcdsRunnerProperty = "tpcds.runner"
 def tpcdsRunnerDependency = project.findProperty(tpcdsRunnerProperty)
         ?: ":runners:direct-java"
 def isSpark = tpcdsRunnerDependency.startsWith(":runners:spark:")
-def isFlink2 = tpcdsRunnerDependency.startsWith(":runners:flink:") && {
+def gtFlink20 = tpcdsRunnerDependency.startsWith(":runners:flink:") && {
   def version = tpcdsRunnerDependency.substring(":runners:flink:".length())
   try {
     def parts = version.split('\\.')
     def major = parts[0].toInteger()
     def minor = parts.length > 1 ? parts[1].toInteger() : 0
-    return (major > 2) || (major == 2 && minor >= 1)
+    return (major == 2 && minor > 0)
   } catch (Exception e) {
     return false
   }
@@ -94,7 +94,7 @@ dependencies {
     gradleRun project(path: tpcdsRunnerDependency, configuration: 
runnerConfiguration)
 }
 
-if (isFlink2) {
+if (!gtFlink20) {
     resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
 }
 
@@ -102,6 +102,7 @@ if (isSpark) {
     configurations.gradleRun {
       exclude group: "org.slf4j", module: "slf4j-jdk14"
     }
+    resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
 }
 
 def sparkJvmArgs() {
diff --git a/sdks/java/testing/watermarks/build.gradle 
b/sdks/java/testing/watermarks/build.gradle
index ca774815467..74955685b7c 100644
--- a/sdks/java/testing/watermarks/build.gradle
+++ b/sdks/java/testing/watermarks/build.gradle
@@ -38,7 +38,18 @@ def runnerProperty = "runner"
 def runnerDependency = (project.hasProperty(runnerProperty)
         ? project.getProperty(runnerProperty)
         : ":runners:direct-java")
-
+def isSparkRunner = runnerDependency.startsWith(":runners:spark:")
+def gtFlink20 = runnerDependency.startsWith(":runners:flink:") && {
+  def version = runnerDependency.substring(":runners:flink:".length())
+  try {
+    def parts = version.split('\\.')
+    def major = parts[0].toInteger()
+    def minor = parts.length > 1 ? parts[1].toInteger() : 0
+    return (major == 2 && minor > 0)
+  } catch (Exception e) {
+    return false
+  }
+}()
 def isDataflowRunner = 
":runners:google-cloud-dataflow-java".equals(runnerDependency)
 def runnerConfiguration = ":runners:direct-java".equals(runnerDependency) ? 
"shadow" : null
 
@@ -74,6 +85,18 @@ dependencies {
   gradleRun project(path: runnerDependency, configuration: runnerConfiguration)
 }
 
+if (!gtFlink20) {
+  resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
+}
+
+if (isSparkRunner) {
+  configurations.gradleRun {
+    // Using Spark runner causes a StackOverflowError if slf4j-jdk14 is on the 
classpath
+    exclude group: "org.slf4j", module: "slf4j-jdk14"
+  }
+  resolveCapabilitiesConflict(configurations.gradleRun, 'org.lz4:lz4-java', 
'at.yawk.lz4')
+}
+
 task run(type: JavaExec) {
   def loadTestArgs = project.findProperty(loadTestArgsProperty) ?: ""
 
diff --git a/settings.gradle.kts b/settings.gradle.kts
index d9dbbf9021e..f71f6de9d16 100644
--- a/settings.gradle.kts
+++ b/settings.gradle.kts
@@ -364,20 +364,8 @@ project(":beam-test-gha").projectDir = file(".github")
 include("beam-validate-runner")
 project(":beam-validate-runner").projectDir = 
file(".test-infra/validate-runner")
 include("com.google.api.gax.batching")
-include("sdks:java:io:kafka:kafka-390")
-findProject(":sdks:java:io:kafka:kafka-390")?.name = "kafka-390"
-include("sdks:java:io:kafka:kafka-312")
-findProject(":sdks:java:io:kafka:kafka-312")?.name = "kafka-312"
-include("sdks:java:io:kafka:kafka-282")
-findProject(":sdks:java:io:kafka:kafka-282")?.name = "kafka-282"
-include("sdks:java:io:kafka:kafka-251")
-findProject(":sdks:java:io:kafka:kafka-251")?.name = "kafka-251"
-include("sdks:java:io:kafka:kafka-241")
-findProject(":sdks:java:io:kafka:kafka-241")?.name = "kafka-241"
-include("sdks:java:io:kafka:kafka-231")
-findProject(":sdks:java:io:kafka:kafka-231")?.name = "kafka-231"
-include("sdks:java:io:kafka:kafka-201")
-findProject(":sdks:java:io:kafka:kafka-201")?.name = "kafka-201"
+include("sdks:java:io:kafka:kafka-392")
+findProject(":sdks:java:io:kafka:kafka-392")?.name = "kafka-392"
 include("sdks:java:managed")
 findProject(":sdks:java:managed")?.name = "managed"
 include("sdks:java:io:iceberg")


Reply via email to