【问题标题】:Kafka consumer for multiple topic多个主题的 Kafka 消费者
【发布时间】:2017-01-26 21:27:18
【问题描述】:

我有一个主题列表(现在是 10 个),它们的大小将来会增加。我知道我们可以从每个主题中生成多个线程(每个主题)来使用,但在我的情况下,如果主题数量增加,那么从主题中消耗的线程数量会增加,这是我不希望的,因为主题不是会过于频繁地获取数据,因此线程将处于理想状态。

有没有办法让一个消费者从所有主题中消费?如果是,那么我们如何才能实现呢?此外,Kafka 将如何维护偏移量?请提出答案。

【问题讨论】:

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


    【解决方案1】:

    我们可以使用以下 API 订阅多个主题: consumer.subscribe(Arrays.asList(topic1,topic2), ConsumerRebalanceListener obj)

    Consumer 有主题信息,我们可以使用 consumer.commitAsync 或 consumer.commitSync() 通过创建 OffsetAndMetadata 对象来提交,如下所示。

    ConsumerRecords<String, String> records = consumer.poll(long value);
    for (TopicPartition partition : records.partitions()) {
        List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);
        for (ConsumerRecord<String, String> record : partitionRecords) {
            System.out.println(record.offset() + ": " + record.value());
        }
        long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
        consumer.commitSync(Collections.singletonMap(partition, new OffsetAndMetadata(lastOffset + 1)));
    }
    

    【讨论】:

    • 我知道,我们可以,但是 Kafka 将如何维护偏移量?另外,有一个单一的消费者群体会解决我的问题吗?
    • 偏移量由您的应用程序提交并存储在名为 __consumer_offsets 的特殊偏移量 kafka 主题中。每个主题的每个分区都会保留偏移量,因此您订阅多少主题并不重要。
    • 这似乎要求所有主题使用相同的序列化程序。有什么方法可以在单个消费者投票中允许每个主题使用不同的序列化程序?
    【解决方案2】:

    不需要多个线程,你可以有一个消费者,从多个主题消费。 偏移量由 zookeeper 维护,因为 kafka-server 本身是无状态的。 每当消费者消费一条消息时,它的偏移量都会提交给zookeeper,以保持未来的跟踪,只处理每条消息一次。所以即使在 kafka 失败的情况下,消费者也会从最后提交的偏移量的下一个开始消费。

    【讨论】:

    • 从 Kafka 0.9 及更高版本开始,偏移量存储在 Kafka 主题而不是 zookeeper 中
    猜你喜欢
    • 1970-01-01
    • 2018-12-31
    • 2018-08-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-10
    • 2020-10-12
    • 2020-06-03
    相关资源
    最近更新 更多