【问题标题】:kafka 0.90 consumer persist group between runs运行之间的 kafka 0.90 消费者持久组
【发布时间】:2016-05-23 15:06:48
【问题描述】:

我已经构建了以下 kafka 消费者:

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:6667");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "TEST1");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "10000");
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,"1000");
this.kconsumer = new KafkaConsumer(props);

我希望消费者在该组启动时最早开始。所以我第一次运行它时,它按预期完美运行。只要订阅存在并且连接没有关闭,它就会继续增加偏移量。

当我登录到 kafka 并运行以下命令时:

./kafka-consumer-groups.sh --bootstrap-server localhost:6667 --new-consumer --group TEST1 --describe

我确切地看到了预期的结果,偏移量的增加等。当连接关闭时,运行相同的命令会导致“消费者组 TEST1 不存在或正在重新平衡。”只是它没有再平衡,它消失了。

当消费者不运行时,如何保持组的存在?我是否缺少消费者或 kafka 中的配置?

另外说明,当我将 OFFSET 参数更改为“最新”时,除非加载新记录,否则即使记录未过期,我也不会得到任何记录。

所以最重要的是,我想做的是启动一个具有给定名称的新消费者,能够从最早的可用记录中提取,关闭该消费者,如果我再次使用该名称启动消费者从我离开的地方拉。关于我所缺少的任何想法?还是我完全误解了高级消费者的工作方式?

【问题讨论】:

  • 我发现如果我将 OFFSET 更改为最新,我会得到想要的结果。但是,我如何检查该组以前是否存在过?因为如果我在从未生成组时将 OFFSET 设置为最新,则它不会返回任何记录。所以看起来我需要一件新的东西,而以前使用过的东西。
  • 您禁用了自动提交 -- 您是否手动提交?如果你不提交,Kafka 在停止消费者时无法知道你离开了哪里。
  • 我手动提交是的。因此,第二次运行消费者时将 OFFSET 更改为最新的工作。

标签: apache-kafka kafka-consumer-api


【解决方案1】:

如果有人遇到这个并想知道我做了什么。在确定组是否首先存在后,我能够设置偏移量。这样做意味着如果该组存在,则使用“最新”。如果不是,请使用“最早”。

    private void buildConsumer(String offset)
    {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:6667");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, this.groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "10000");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, offset);
        props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,"1000");
        this.kconsumer = new KafkaConsumer(props);
    }

    /*
    Check if the group exists before polling.
    If it does, leave with default offset.
    If it does not exists, set the offset to earliest to ensure you are getting all the records
    */
    private void groupExists(String topic)
    {
        TopicPartition toc = new TopicPartition(topic, 0);
        OffsetAndMetadata oam = kconsumer.committed(toc);
        if(oam != null){
            //do nothing, all is well, start from last commit
        } else {
            /*
            when a new group is started the AUTO_OFFSET_RESET_CONFIG
            needs to be set to earliest to ensure all records are picked up
            Since that property can only be set at instantiation the consumer
            must be rebuilt and resubscribed
            */
            buildConsumer("earliest");
            this.kconsumer.subscribe(Arrays.asList(topic));
        }
    }

【讨论】:

    猜你喜欢
    • 2016-06-02
    • 2016-10-31
    • 1970-01-01
    • 2020-09-19
    • 1970-01-01
    • 2017-01-04
    • 2020-06-11
    • 1970-01-01
    • 2022-11-09
    相关资源
    最近更新 更多