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")