Skip to content

Commit 44bd0ef

Browse files
committed
Replace obsolete class KafkaServerStartable with KafkaServer
1 parent 6e42c7d commit 44bd0ef

4 files changed

Lines changed: 21 additions & 9 deletions

File tree

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -842,6 +842,7 @@ class BeamModulePlugin implements Plugin<Project> {
842842
kafka_scala_2_12 : "org.apache.kafka:kafka_2.12:$kafka_version",
843843
kafka_scala_2_13 : "org.apache.kafka:kafka_2.13:$kafka_version",
844844
kafka_clients : "org.apache.kafka:kafka-clients:$kafka_version",
845+
kafka_server : "org.apache.kafka:kafka-server:$kafka_version",
845846
log4j : "log4j:log4j:1.2.17",
846847
log4j_over_slf4j : "org.slf4j:log4j-over-slf4j:$slf4j_version",
847848
log4j2_api : "org.apache.logging.log4j:log4j-api:$log4j2_version",

runners/spark/spark_runner.gradle

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -289,6 +289,7 @@ dependencies {
289289
testImplementation library.java.avro
290290
testImplementation (spark_scala_version == '2.13' ? library.java.kafka_scala_2_13 : library.java.kafka_scala_2_12)
291291
testImplementation library.java.kafka_clients
292+
testImplementation library.java.kafka_server
292293
testImplementation library.java.junit
293294
testImplementation library.java.mockito_core
294295
testImplementation "org.assertj:assertj-core:3.11.1"

runners/spark/src/test/java/org/apache/beam/runners/spark/translation/streaming/utils/EmbeddedKafkaCluster.java

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@
2929
import java.util.Properties;
3030
import java.util.Random;
3131
import kafka.server.KafkaConfig;
32-
import kafka.server.KafkaServerStartable;
32+
import kafka.server.KafkaServer;
3333
import org.apache.zookeeper.server.NIOServerCnxnFactory;
3434
import org.apache.zookeeper.server.ServerCnxnFactory;
3535
import org.apache.zookeeper.server.ZooKeeperServer;
@@ -47,7 +47,7 @@ public class EmbeddedKafkaCluster {
4747

4848
private final String brokerList;
4949

50-
private final List<KafkaServerStartable> brokers;
50+
private final List<KafkaServer> brokers;
5151
private final List<File> logDirs;
5252

5353
private EmbeddedKafkaCluster(String zkConnection) {
@@ -114,15 +114,20 @@ public void startup() {
114114
properties.setProperty("offsets.topic.replication.factor", "1");
115115
properties.setProperty("log.flush.interval.messages", String.valueOf(1));
116116

117-
KafkaServerStartable broker = startBroker(properties);
117+
KafkaServer broker = startBroker(properties);
118118

119119
brokers.add(broker);
120120
logDirs.add(logDir);
121121
}
122122
}
123123

124-
private static KafkaServerStartable startBroker(Properties props) {
125-
KafkaServerStartable server = new KafkaServerStartable(new KafkaConfig(props));
124+
private static KafkaServer startBroker(Properties props) {
125+
KafkaServer server =
126+
new KafkaServer(
127+
new KafkaConfig(props),
128+
KafkaServer.$lessinit$greater$default$2(),
129+
KafkaServer.$lessinit$greater$default$3(),
130+
KafkaServer.$lessinit$greater$default$4());
126131
server.startup();
127132
return server;
128133
}
@@ -148,7 +153,7 @@ public String getZkConnection() {
148153

149154
@SuppressWarnings("Slf4jDoNotLogMessageOfExceptionExplicitly")
150155
public void shutdown() {
151-
for (KafkaServerStartable broker : brokers) {
156+
for (KafkaServer broker : brokers) {
152157
try {
153158
broker.shutdown();
154159
} catch (Exception e) {

sdks/java/testing/kafka-service/src/test/java/org/apache/beam/sdk/testing/kafka/LocalKafka.java

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,18 +20,23 @@
2020
import java.nio.file.Files;
2121
import java.util.Properties;
2222
import kafka.server.KafkaConfig;
23-
import kafka.server.KafkaServerStartable;
23+
import kafka.server.KafkaServer;
2424

2525
public class LocalKafka {
26-
private final KafkaServerStartable server;
26+
private final KafkaServer server;
2727

2828
LocalKafka(int kafkaPort, int zookeeperPort) throws Exception {
2929
Properties kafkaProperties = new Properties();
3030
kafkaProperties.setProperty("port", String.valueOf(kafkaPort));
3131
kafkaProperties.setProperty("zookeeper.connect", String.format("localhost:%s", zookeeperPort));
3232
kafkaProperties.setProperty("offsets.topic.replication.factor", "1");
3333
kafkaProperties.setProperty("log.dir", Files.createTempDirectory("kafka-log-").toString());
34-
server = new KafkaServerStartable(KafkaConfig.fromProps(kafkaProperties));
34+
server =
35+
new KafkaServer(
36+
KafkaConfig.fromProps(kafkaProperties),
37+
KafkaServer.$lessinit$greater$default$2(),
38+
KafkaServer.$lessinit$greater$default$3(),
39+
KafkaServer.$lessinit$greater$default$4());
3540
}
3641

3742
public void start() {

0 commit comments

Comments
 (0)