【问题标题】:Read unprocessed messages in Apache Kafka using Simple Consumer使用 Simple Consumer 读取 Apache Kafka 中未处理的消息
【发布时间】:2015-02-03 08:47:27
【问题描述】:

我厌倦了点击链接

https://cwiki.apache.org/confluence/display/KAFKA/0.8.0+SimpleConsumer+Example

使用 SimpleConsumer 来消费消息,但是在使用它时我发现了一些突然的行为,如下所示:

消费者正在消费来自特定分区的消息。但问题是,当我的消费者正在运行并且我使用生产者将消息推送到主题时,它会使用来自该分区的消息。但是,如果我的消费者目前没有运行并且我将一些消息推送到主题并再次启动消费者,它不会消费生产者推送的消息,但它再次准备好消费现在将被推送的消息。我正在使用 LatestTime() 而不是 EarliestTime() 因为我只想使用未处理的消息。

例如

案例-1

消费者正在运行:

Producer将M1、M2、M3消息推送到topic 1的partition 1

结果:消费者将消费所有三个消息。

案例 - 2

消费者没有运行

producer 现在将 m4、m5 m6 messgae 推送到主题 1 的分区 1

现在调用消费者

结果:消费者不消费消息 m4、m5、m6 但如果我要检查偏移量,则它设置为 7。这意味着生产者在生成消息时已将偏移量提前到 7,因此消费者将使用来自的消息现在偏移 7

当消费者再次出现时,它应该从 m4 读取消息时,请提供帮助。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    你做错了。

    首先,我不确定SimpleConsumer 是您要查找的内容。它迫使您自己管理偏移量(例如,它根本不会向 Zookeeper 提交偏移量,并且每次您再次启动 SimpleConsumer 时,它都会再次获取相同的消息)。 SimpleConsumer 不了解“已处理消息”。它所能做的就是从某个偏移量开始获取并继续获取,直到你说“停止”。

    无论如何,如果您打算自己提交已处理的偏移量,您应该使用EarliestTimeauto.offset.reset=smallest 配置条目)。 auto.offset.reset 表示如果您的消费者使用错误的偏移量初始化(如果我没记错,SimpleConsumer 使用-1 偏移量初始化,这显然是错误的)它将重置为smallest 可用(EarliestTime)或largest 可用 (LatestTime) 偏移量。

    为了更清楚,这里是一个例子:

    你的Case-1:

    您创建一个消费者并将其指向主题 1 分区 1。由于它最初使用错误的偏移量进行初始化,它会要求代理提供一些适当的偏移量(这里是 smallestlargest 偏移量重置的来源)。如果您还没有产生任何消息,smallestlargest 偏移量都将是 0,因此当您产生一些消息时,您的消费者将获取这些消息。

    Case-2:

    你产生 N 条消息(比如 7 条)。然后你开始你的SimpleConsumer。同样,它使用错误的偏移量进行初始化,并要求代理提供正确的偏移量。使用smallest 重置偏移量将是0,使用largest 偏移量它将是7。在您的示例中,您使用 LargestOffsets 您的消费者将使用偏移量 7 重新初始化并开始使用它。

    一般来说,看看高级消费者,在大多数情况下,这就是您要寻找的东西。 这是link

    【讨论】:

    • 但如果我使用 HighLevelConsumer,那么它不会为我提供将消费者与 toipc 中的特定分区相关联的方法。我的要求是 - 我有 n 个生产者,每个生产者都附加到一个特定的分区,n 个消费者每个附加到一个特定的分区。分区是根据一些键完成的。所以我在这种情况下选择了 SimpleConsumer。
    • @ajaymittal 那么不幸的是,您将不得不自己管理偏移量。高级消费者不提供消费特定分区的配置。
    • 在使用 SimpleConsumer 时,我只能使用 EarliestTime() 作为 0 获得最小的偏移量,或者 latestTime() 给出最近的偏移量。有没有办法让偏移点指向未处理消息的索引。
    • 您必须以某种方式自己跟踪已处理的偏移量。高级消费者为此使用 Zookeeper。 SimpleConsumer 从不提交消耗的偏移量,因此它不知道已经处理了哪些消息。您必须以某种方式自己跟踪最后处理的偏移量,当您重新启动消费者时,您应该使用此存储的偏移量手动对其进行初始化。
    猜你喜欢
    • 2013-04-24
    • 1970-01-01
    • 2019-01-23
    • 2022-01-06
    • 1970-01-01
    • 2016-06-29
    • 1970-01-01
    • 2018-07-30
    • 2022-10-14
    相关资源
    最近更新 更多