You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Docker Kafka容器消费者无法消费数据问题求助

Fixing Java Producer/Consumer Issues with spotify/kafka Docker Container

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 set ADVERTISED_HOST and ADVERTISED_PORT correctly when starting the container, and that your bootstrap servers match that address.
  • NoBrokerAvailableException: Verify the container is running (docker ps to 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 10:12:49