advertisedListeners = new ArrayList<>();
advertisedListeners.add("PLAINTEXT://" + getBootstrapServers());
advertisedListeners.add(brokerAdvertisedListener);
+
+ advertisedListeners.addAll(KafkaHelper.resolveAdvertisedListeners(this.advertisedListeners));
String kafkaAdvertisedListeners = String.join(",", advertisedListeners);
+
String command = "#!/bin/bash\n";
// exporting KAFKA_ADVERTISED_LISTENERS with the container hostname
command += String.format("export KAFKA_ADVERTISED_LISTENERS=%s\n", kafkaAdvertisedListeners);
@@ -89,6 +80,69 @@ protected void containerIsStarting(InspectContainerResponse containerInfo) {
copyFileToContainer(Transferable.of(command, 0777), STARTER_SCRIPT);
}
+ /**
+ * Add a listener in the format {@code host:port}.
+ * Host will be included as a network alias.
+ *
+ * Use it to register additional connections to the Kafka broker within the same container network.
+ *
+ * The listener will be added to the list of default listeners.
+ *
+ * Default listeners:
+ *
+ * - 0.0.0.0:9092
+ * - 0.0.0.0:9093
+ * - 0.0.0.0:9094
+ *
+ *
+ * The listener will be added to the list of default advertised listeners.
+ *
+ * Default advertised listeners:
+ *
+ * - {@code container.getConfig().getHostName():9092}
+ * - {@code container.getHost():container.getMappedPort(9093)}
+ *
+ * @param listener a listener with format {@code host:port}
+ * @return this {@link KafkaContainer} instance
+ */
+ public KafkaContainer withListener(String listener) {
+ this.listeners.add(listener);
+ this.advertisedListeners.add(() -> listener);
+ return this;
+ }
+
+ /**
+ * Add a listener in the format {@code host:port} and a {@link Supplier} for the advertised listener.
+ * Host from listener will be included as a network alias.
+ *
+ * Use it to register additional connections to the Kafka broker from outside the container network
+ *
+ * The listener will be added to the list of default listeners.
+ *
+ * Default listeners:
+ *
+ * - 0.0.0.0:9092
+ * - 0.0.0.0:9093
+ * - 0.0.0.0:9094
+ *
+ *
+ * The {@link Supplier} will be added to the list of default advertised listeners.
+ *
+ * Default advertised listeners:
+ *
+ * - {@code container.getConfig().getHostName():9092}
+ * - {@code container.getHost():container.getMappedPort(9093)}
+ *
+ * @param listener a supplier that will provide a listener
+ * @param advertisedListener a supplier that will provide a listener
+ * @return this {@link KafkaContainer} instance
+ */
+ public KafkaContainer withListener(String listener, Supplier advertisedListener) {
+ this.listeners.add(listener);
+ this.advertisedListeners.add(advertisedListener);
+ return this;
+ }
+
public String getBootstrapServers() {
return String.format("%s:%s", getHost(), getMappedPort(KAFKA_PORT));
}
diff --git a/modules/kafka/src/test/java/org/testcontainers/kafka/CompatibleApacheKafkaImageTest.java b/modules/kafka/src/test/java/org/testcontainers/kafka/CompatibleApacheKafkaImageTest.java
new file mode 100644
index 00000000000..b2bb31ae502
--- /dev/null
+++ b/modules/kafka/src/test/java/org/testcontainers/kafka/CompatibleApacheKafkaImageTest.java
@@ -0,0 +1,26 @@
+package org.testcontainers.kafka;
+
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.Parameterized;
+import org.testcontainers.AbstractKafka;
+
+@RunWith(Parameterized.class)
+public class CompatibleApacheKafkaImageTest extends AbstractKafka {
+
+ @Parameterized.Parameters(name = "{0}")
+ public static String[] params() {
+ return new String[] { "apache/kafka:3.8.0", "apache/kafka-native:3.8.0" };
+ }
+
+ @Parameterized.Parameter
+ public String imageName;
+
+ @Test
+ public void testUsage() throws Exception {
+ try (KafkaContainer kafka = new KafkaContainer(this.imageName)) {
+ kafka.start();
+ testKafkaFunctionality(kafka.getBootstrapServers());
+ }
+ }
+}
diff --git a/modules/kafka/src/test/java/org/testcontainers/kafka/KafkaContainerTest.java b/modules/kafka/src/test/java/org/testcontainers/kafka/KafkaContainerTest.java
index 233ee1cd2d4..b1764d3b522 100644
--- a/modules/kafka/src/test/java/org/testcontainers/kafka/KafkaContainerTest.java
+++ b/modules/kafka/src/test/java/org/testcontainers/kafka/KafkaContainerTest.java
@@ -1,26 +1,61 @@
package org.testcontainers.kafka;
import org.junit.Test;
-import org.junit.runner.RunWith;
-import org.junit.runners.Parameterized;
import org.testcontainers.AbstractKafka;
+import org.testcontainers.KCatContainer;
+import org.testcontainers.containers.Network;
+import org.testcontainers.containers.SocatContainer;
+
+import static org.assertj.core.api.Assertions.assertThat;
-@RunWith(Parameterized.class)
public class KafkaContainerTest extends AbstractKafka {
- @Parameterized.Parameters(name = "{0}")
- public static String[] params() {
- return new String[] { "apache/kafka:3.8.0", "apache/kafka-native:3.8.0" };
+ @Test
+ public void testUsage() throws Exception {
+ try ( // constructorWithVersion {
+ KafkaContainer kafka = new KafkaContainer("apache/kafka-native:3.8.0")
+ // }
+ ) {
+ kafka.start();
+ testKafkaFunctionality(kafka.getBootstrapServers());
+ }
}
- @Parameterized.Parameter
- public String imageName;
+ @Test
+ public void testUsageWithListener() throws Exception {
+ try (
+ Network network = Network.newNetwork();
+ KafkaContainer kafka = new KafkaContainer("apache/kafka-native:3.8.0")
+ .withListener("kafka:19092")
+ .withNetwork(network);
+ KCatContainer kcat = new KCatContainer().withNetwork(network)
+ ) {
+ kafka.start();
+ kcat.start();
+
+ kcat.execInContainer("kcat", "-b", "kafka:19092", "-t", "msgs", "-P", "-l", "/data/msgs.txt");
+ String stdout = kcat
+ .execInContainer("kcat", "-b", "kafka:19092", "-C", "-t", "msgs", "-c", "1")
+ .getStdout();
+
+ assertThat(stdout).contains("Message produced by kcat");
+ }
+ }
@Test
- public void testUsage() throws Exception {
- try (KafkaContainer kafka = new KafkaContainer(imageName)) {
+ public void testUsageWithListenerFromProxy() throws Exception {
+ try (
+ Network network = Network.newNetwork();
+ SocatContainer socat = new SocatContainer().withNetwork(network).withTarget(2000, "kafka", 19092);
+ KafkaContainer kafka = new KafkaContainer("apache/kafka-native:3.8.0")
+ .withListener("kafka:19092", () -> socat.getHost() + ":" + socat.getMappedPort(2000))
+ .withNetwork(network)
+ ) {
+ socat.start();
kafka.start();
- testKafkaFunctionality(kafka.getBootstrapServers());
+
+ String bootstrapServers = String.format("%s:%s", socat.getHost(), socat.getMappedPort(2000));
+ testKafkaFunctionality(bootstrapServers);
}
}
}