【问题标题】:Kafka last message poll results in 0 messagesKafka 上次消息轮询结果为 0 条消息
【发布时间】:2019-01-09 10:06:37
【问题描述】:

我有一个带有单个分区的 Kafka 主题 (1.0.0)。消费者被打包在一个 EAR 中,当部署到 Wildfly 10 时,最后一条消息的轮询总是返回 0 条消息。虽然题目不是空的。

final TopicPartition tp = new TopicPartition(topic, 0);

final Long beginningOffset = consumer.beginningOffsets(Collections.singleton(tp)).get(tp);
final Long endOffset = consumer.endOffsets(Collections.singleton(tp)).get(tp);

consumer.assign(Collections.singleton(tp));
consumer.seek(tp, endOffset - 1); 

当我进行投票时,我得到 0 条记录。尽管记录表明:

Consumer is now at position 377408 while Topic begin is 0 and end is 377409

当我更改为 -2 时:

consumer.seek(tp, endOffset - 2);

我确实收到一条消息:

 Consumer is now at position 377407 while Topic begin is 0 and end is 377409

但这当然不是正确的记录,消息 377408 在哪里?

尝试了很多方法来寻求结束等,但它从来没有奏效。

这是我的消费者配置:

Properties properties = new Properties();
properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, Configuration.KAFKA_SERVERS.getAsString());
properties.put(ConsumerConfig.GROUP_ID_CONFIG, GROUP_ID);
properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class);
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
properties.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");

注意:我尝试了 read_uncommitted 和 read_committed,都给出了相同的结果。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    正如Javadoc中提到的,这是因为endOffsets()返回:

    最后一次成功复制消息的偏移量加一

    这实际上是下一条消息将获得的偏移量。

    这就是为什么寻找endOffset - 1 不会返回任何内容,而寻找endOffset - 2 只会返回最后一条消息。

    我同意这可能不是最直观的行为,但这就是它目前的工作方式!

    【讨论】:

    • 不,看调试语句,结尾是 377409 而消费者是 377408,所以 endoffset 减一应该会导致消息 377408。我的代码也是基于几个关于如何获取的示例最后一条消息。
    • 另外,在查看 Kafka 文档时:consumer.position --> 获取将要获取的 下一条记录 的偏移量(如果存在具有该偏移量的记录)。那么那个记录在哪里呢?
    • 经过更多测试,我会接受这个作为答案。尽管这些文档似乎违反直觉。
    猜你喜欢
    • 2019-09-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-29
    • 2016-12-19
    • 1970-01-01
    • 2017-08-30
    • 2021-09-09
    相关资源
    最近更新 更多