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 dc88f3fe0a5a..8c61cb4e15f2 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 @@ -27,7 +27,7 @@ public class LocalKafka { LocalKafka(int kafkaPort, int zookeeperPort) throws Exception { Properties kafkaProperties = new Properties(); - kafkaProperties.setProperty("port", String.valueOf(kafkaPort)); + kafkaProperties.setProperty("listeners", String.format("PLAINTEXT://localhost:%s", kafkaPort)); 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());