Skip to content

Commit ec397ce

Browse files
authored
[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
1 parent 804b127 commit ec397ce

37 files changed

Lines changed: 167 additions & 237 deletions

File tree

buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -639,7 +639,7 @@ class BeamModulePlugin implements Plugin<Project> {
639639
def jaxb_api_version = "2.3.3"
640640
def jsr305_version = "3.0.2"
641641
def everit_json_version = "1.14.2"
642-
def kafka_version = "2.4.1"
642+
def kafka_version = "3.9.2"
643643
def log4j2_version = "2.25.4"
644644
def nemo_version = "0.1"
645645
// [bomupgrader] determined by: io.grpc:grpc-netty, consistent with: google_cloud_platform_libraries_bom
@@ -850,8 +850,10 @@ class BeamModulePlugin implements Plugin<Project> {
850850
jupiter_api : "org.junit.jupiter:junit-jupiter-api:$jupiter_version",
851851
jupiter_engine : "org.junit.jupiter:junit-jupiter-engine:$jupiter_version",
852852
jupiter_params : "org.junit.jupiter:junit-jupiter-params:$jupiter_version",
853-
kafka : "org.apache.kafka:kafka_2.11:$kafka_version",
853+
kafka_scala_2_12 : "org.apache.kafka:kafka_2.12:$kafka_version",
854+
kafka_scala_2_13 : "org.apache.kafka:kafka_2.13:$kafka_version",
854855
kafka_clients : "org.apache.kafka:kafka-clients:$kafka_version",
856+
kafka_server : "org.apache.kafka:kafka-server:$kafka_version",
855857
log4j : "log4j:log4j:1.2.17",
856858
log4j_over_slf4j : "org.slf4j:log4j-over-slf4j:$slf4j_version",
857859
log4j2_api : "org.apache.logging.log4j:log4j-api:$log4j2_version",

examples/java/build.gradle

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,8 +44,8 @@ dependencies {
4444
if (project.findProperty('testJavaVersion') == '21' || JavaVersion.current().compareTo(JavaVersion.VERSION_21) >= 0) {
4545
// this dependency is a provided dependency for kafka-avro-serializer. It is not needed to compile with Java<=17
4646
// but needed for compile only under Java21, specifically, required for extending from AbstractKafkaAvroDeserializer
47-
compileOnly library.java.kafka
48-
permitUnusedDeclared library.java.kafka
47+
compileOnly library.java.kafka_scala_2_12
48+
permitUnusedDeclared library.java.kafka_scala_2_12
4949
}
5050
implementation library.java.kafka_clients
5151
implementation project(path: ":sdks:java:core", configuration: "shadow")

examples/java/common.gradle

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ configurations.sparkRunnerPreCommit {
3636
exclude group: "org.slf4j", module: "slf4j-jdk14"
3737
}
3838
resolveCapabilitiesConflict(configurations.flinkRunnerPreCommit, 'org.lz4:lz4-java', 'at.yawk.lz4')
39+
resolveCapabilitiesConflict(configurations.sparkRunnerPreCommit, 'org.lz4:lz4-java', 'at.yawk.lz4')
3940

4041

4142
dependencies {

it/kafka/src/test/java/org/apache/beam/it/kafka/KafkaResourceManagerTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -166,7 +166,8 @@ public void testCleanupShouldDropNonStaticTopic() throws IOException {
166166
KafkaResourceManager tm = new KafkaResourceManager(kafkaClient, container, builder);
167167

168168
tm.cleanupAll();
169-
verify(kafkaClient).deleteTopics(argThat(list -> list.size() == numTopics));
169+
verify(kafkaClient)
170+
.deleteTopics(argThat((Collection<String> list) -> list.size() == numTopics));
170171
}
171172

172173
@Test

runners/flink/1.19/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,3 +23,9 @@ project.ext {
2323

2424
// Load the main build script which contains all build logic.
2525
apply from: "../flink_runner.gradle"
26+
27+
// Flink 1.19 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
28+
// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
29+
configurations.all {
30+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
31+
}

runners/flink/1.19/job-server/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,3 +29,9 @@ project.ext {
2929

3030
// Load the main build script which contains all build logic.
3131
apply from: "$basePath/flink_job_server.gradle"
32+
33+
// Flink 1.19 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
34+
// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
35+
configurations.all {
36+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
37+
}

runners/flink/1.20/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,3 +23,9 @@ project.ext {
2323

2424
// Load the main build script which contains all build logic.
2525
apply from: "../flink_runner.gradle"
26+
27+
// Flink 1.20 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
28+
// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
29+
configurations.all {
30+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
31+
}

runners/flink/1.20/job-server/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,3 +29,9 @@ project.ext {
2929

3030
// Load the main build script which contains all build logic.
3131
apply from: "$basePath/flink_job_server.gradle"
32+
33+
// Flink 1.20 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
34+
// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
35+
configurations.all {
36+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
37+
}

runners/flink/2.0/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,3 +41,9 @@ project.ext {
4141

4242
// Load the main build script which contains all build logic.
4343
apply from: "../flink_runner.gradle"
44+
45+
// Flink 2.0 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
46+
// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
47+
configurations.all {
48+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
49+
}

runners/flink/2.0/job-server/build.gradle

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,3 +29,9 @@ project.ext {
2929

3030
// Load the main build script which contains all build logic.
3131
apply from: "$basePath/flink_job_server.gradle"
32+
33+
// Flink 2.0 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
34+
// Explicitly prefer at.yawk.lz4 candidates to resolve capability conflict
35+
configurations.all {
36+
resolveCapabilitiesConflict(it, 'org.lz4:lz4-java', 'at.yawk.lz4')
37+
}

0 commit comments

Comments
 (0)