【发布时间】: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