【问题标题】:Kafka High Level Consumer Fetch All Messages From Topic Using Java API (Equivalent to --from-beginning)Kafka 高级消费者使用 Java API 从主题中获取所有消息(相当于 --from-beginning)
【发布时间】:2014-03-12 19:31:18
【问题描述】:

我正在使用来自 Kafka 站点的 ConsumerGroupExample 代码测试 Kafka 高级消费者。我想检索 Kafka 服务器配置中名为“test”的主题的所有现有消息。查看其他博客,auto.offset.reset 应该设置为“最小”才能获取所有消息:

private static ConsumerConfig createConsumerConfig(String a_zookeeper, String a_groupId)    {
    Properties props = new Properties();
    props.put("zookeeper.connect", a_zookeeper);
    props.put("group.id", a_groupId);
    props.put("auto.offset.reset", "smallest");
    props.put("zookeeper.session.timeout.ms", "10000");     

    return new ConsumerConfig(props);
}

我真正遇到的问题是:高级消费者的等效 Java api 调用是什么,相当于:

bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic test --from-beginning

【问题讨论】:

    标签: apache-kafka java consumer


    【解决方案1】:

    看来您需要使用“低级 SimpleConsumer API”

    对于大多数应用程序,高级消费者 Api 已经足够好了。 一些应用程序需要不向高级消费者公开的功能 然而(例如,在重新启动消费者时设置初始偏移量)。他们能 而是使用我们的低级 SimpleConsumer Api。逻辑会有点 比较复杂的可以参考here中的例子。

    此示例用于从具有以下参数的主题中获取所有消息:(请注意,端口是 Kafka 端口,而不是 ZooKeeper 端口,主题设置自 this example):

    10 my-replicated-topic 0 localhost 9092
    

    具体来说,有一种获取readOffset的方法,它采用kafka.api.OffsetRequest.EarliestTime():

    long readOffset = getLastOffset(consumer,a_topic, a_partition, kafka.api.OffsetRequest.EarliestTime(), clientName);
    

    这里有另一篇文章可能会提供一些关于如何解决这个问题的替代想法:How to get data from old offset point in Kafka?

    【讨论】:

    • 你用什么来实现的?从主题中读取所有消息。
    【解决方案2】:

    基本上,每次新消费者尝试消费主题时,它都会从头开始读取消息。如果您每次都特别从头开始消费以进行测试,那么每次您使用新的 groupID 初始化消费者时,它都会从头开始读取消息。我是这样做的:

    properties.put("group.id", UUID.randomUUID().toString());
    

    每次都从头开始阅读消息!

    【讨论】:

    • 谢谢!需要它用于测试目的。我想这是因为您可以将相同的数据用于不同的目的?
    • @user1758777 是的,我需要每个不同的测试来处理相同的数据。
    【解决方案3】:

    要从头获取消息,您可以这样做:

    import kafka.utils.ZkUtils;
    ZkUtils.maybeDeletePath("zkhost:zkport", "/consumers/group.id");
    

    那就按照常规工作吧……

    【讨论】:

      【解决方案4】:
       Properties props = new Properties(); 
       props.put("bootstrap.servers", "localhost:9092");
       props.put("auto.offset.reset", "earliest");
       props.put("group.id", UUID.randomUUID().toString());
      

      这些属性会帮助你。

      【讨论】:

        猜你喜欢
        • 2016-01-07
        • 1970-01-01
        • 2018-07-27
        • 2016-08-27
        • 1970-01-01
        • 2016-03-11
        • 1970-01-01
        • 2018-12-15
        • 2019-06-25
        相关资源
        最近更新 更多