【发布时间】:2015-08-11 02:03:52
【问题描述】:
我正在使用火花流从 Kafka 读取消息,它工作正常。但是我有一个需要重新阅读消息的要求。我在想我可能只需要更改 spark 的客户 groupId 并重新启动 spark 流应用程序,它应该从头开始重新读取 kafka 消息。但结果是 Spark 无法收到任何消息,我很困惑。通过 Kafka 文档,如果您更改客户 groupId,那么它应该从一开始就收到消息,因为 kafka 将您视为新客户。提前致谢!
【问题讨论】:
我正在使用火花流从 Kafka 读取消息,它工作正常。但是我有一个需要重新阅读消息的要求。我在想我可能只需要更改 spark 的客户 groupId 并重新启动 spark 流应用程序,它应该从头开始重新读取 kafka 消息。但结果是 Spark 无法收到任何消息,我很困惑。通过 Kafka 文档,如果您更改客户 groupId,那么它应该从一开始就收到消息,因为 kafka 将您视为新客户。提前致谢!
【问题讨论】:
Kafka 消费者有一个名为 auto.offset.reset 的属性(参见Kafka Doc)。这告诉消费者在开始消费但尚未提交偏移量时要做什么。这是你的情况。该主题有消息,但没有存储起始偏移量,因为您还没有阅读该新组 ID 下的任何内容。在这种情况下,使用 auto.offset.reset 属性。如果值为“最大”,这是默认值),则开始位置设置为最大偏移量(最后一个),您将获得您所看到的行为。如果该值为“最小”,则偏移量设置为开始偏移量,消费者将读取整个分区。这就是你想要的。
所以我不确定你是如何在 Spark 应用程序中设置 Kafka 属性的,但如果你希望新的组 id 导致阅读整个主题,你肯定希望该属性设置为“最小” .
【讨论】:
听起来您正在为 Kafka 使用 Spark Streaming 基于接收器的 API。正如您所注意到的,对于该 api,auto.offset.reset 仅适用于 ZK 中没有偏移量的情况。
如果您希望能够指定确切的偏移量,请查看将 fromOffsets 作为参数的 createDirectStream 调用版本。
【讨论】: