Docker Kafka容器消费者无法消费数据问题求助
Hey there! I totally get the confusion when switching from regular Kafka to containerized Kafka—those tiny configuration gaps can trip you up fast. Let’s break down what’s going wrong and give you a working example tailored specifically for the spotify/kafka container.
Key Difference: Kafka Broker Advertisement
The biggest gotcha with containerized Kafka vs. regular Kafka is how the broker tells clients where to connect. By default, the spotify/kafka container uses its internal hostname for advertising, which your host machine can’t reach. That’s exactly why your regular Java code throws exceptions—your client is trying to connect to a non-existent address outside the container.
First, let’s start the container correctly with advertised listeners set to your local machine:
docker run -d -p 2181:2181 -p 9092:9092 \ --env ADVERTISED_HOST=localhost \ --env ADVERTISED_PORT=9092 \ spotify/kafka
This tells Kafka to advertise itself as localhost:9092, so your Java client can reach it from your host machine without issues.
Working Java Producer Example
First, add the Kafka clients dependency to your pom.xml (if using Maven):
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.6.1</version> <!-- Use the latest compatible version --> </dependency>
Here’s a simple producer that plays nicely with the container:
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class ContainerKafkaProducer { public static void main(String[] args) { Properties props = new Properties(); // Critical: Match the advertised address we set when starting the container props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) { String topic = "test-topic"; for (int i = 0; i < 5; i++) { ProducerRecord<String, String> record = new ProducerRecord<>(topic, "key-" + i, "message-" + i); producer.send(record, (metadata, exception) -> { if (exception != null) { System.err.println("Oops, error sending message: " + exception.getMessage()); } else { System.out.printf("Message sent to topic %s, partition %d, offset %d%n", metadata.topic(), metadata.partition(), metadata.offset()); } }); } producer.flush(); } } }
Working Java Consumer Example
And here’s a matching consumer to receive those messages:
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class ContainerKafkaConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // Start from the first message if no offset exists try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) { String topic = "test-topic"; consumer.subscribe(Collections.singletonList(topic)); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { System.out.printf("Got message: key=%s, value=%s, partition=%d, offset=%d%n", record.key(), record.value(), record.partition(), record.offset()); } } } } }
Troubleshooting Common Exceptions
TimeoutException/LeaderNotAvailableException: This almost always means your client can’t reach the broker. Double-check that you setADVERTISED_HOSTandADVERTISED_PORTcorrectly when starting the container, and that your bootstrap servers match that address.NoBrokerAvailableException: Verify the container is running (docker psto check) and that port 9092 is properly mapped to your host machine.
Run the producer first (it’ll auto-create the topic since spotify/kafka has auto-topic creation enabled by default), then start the consumer—you should see messages flowing smoothly!
内容的提问来源于stack exchange,提问作者JollyRoger

