【问题标题】:Inconsistent data output from Kafka consumer来自 Kafka 消费者的数据输出不一致
【发布时间】:2018-04-24 16:42:31
【问题描述】:

我需要从 Kafka 消费者那里提取数据以将其传递给我的应用程序。以下是我为访问消费者而编写的代码:

public class ConsumerGroup {
    public static void main(String[] args) throws Exception {

        String topic = "kafka_topic";
        String group = "0";
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", group);
        props.put("enable.auto.commit", "true");
        props.put("auto.commit.interval.ms", "1000");
        props.put("session.timeout.ms", "30000");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("auto.offset.reset", "earliest");
        KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(props);

        consumer.subscribe(Arrays.asList(topic));
        System.out.println("Subscribed to topic: " + topic);

        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(100);
            for (ConsumerRecord<String, String> record : records)
                System.out.printf("offset = %d, key = %s, value = %s\n", record.offset(), record.key(), record.value());
        }
    }
}

当我运行此代码时,有时会生成数据,有时不会生成数据。为什么这种行为不一致?我的代码有问题吗?

【问题讨论】:

    标签: java apache-kafka kafka-consumer-api


    【解决方案1】:

    您的代码没问题。您启用了自动提交选项,因此在您阅读记录后,它们会自动提交给 Kafka。每次运行代码时,您都会从上次处理的偏移量开始,该偏移量存储在 __consumer_offsets 主题中。所以你总是只读取上次运行后到达 Kafka 的新记录。要在消费者应用中不断打印数据,您应该不断将新记录放入您的主题中。

    【讨论】:

    • 感谢您提供此信息。我不想只消耗上次命中后插入的最新记录。我想获取 kafka 中存在的所有数据,因此在获取后我不会删除 kafka 主题中的数据。有没有办法获取所有数据?
    • 您应该使用KafkaConsumer.seekToBeginning 方法,或者每次使用 group.id 属性的新值作为每个 group.id 存储的偏移量
    • 我只有一个主题并且我没有使用任何分区,在调用方法 KafkaConsumer.seekToBeginnng(Collection partitions) 时应该在参数“Collection”中传递什么值?
    • Set&lt;TopicPartition&gt; assignments = consumer.assignment(); assignments.forEach(topicPartition -&gt; consumer.seekToBeginning(Arrays.asList(topicPartition)));
    猜你喜欢
    • 2017-04-14
    • 2019-07-01
    • 2019-08-10
    • 2015-07-24
    • 1970-01-01
    • 2018-08-04
    • 2023-02-02
    • 2021-07-09
    • 2020-07-27
    相关资源
    最近更新 更多