【问题标题】:kafka offset in spark火花中的kafka偏移
【发布时间】:2021-11-16 16:42:44
【问题描述】:

Kafka enable.auto.commit 设置为 false,Spark 版本为 2.4

  1. 如果使用最新的偏移量,我们是否需要手动查找最后的偏移量详细信息并在 Spark 应用程序的 .CreateDirectStream() 中提及?还是会自动采用最新的偏移量?无论如何,我们是否需要手动查找最后的偏移量详细信息。

  2. 使用SparkSession.readstrem.format(kafka)....KafkaUtils.createDirectStream()有什么区别吗?

  3. 使用最早偏移选项时,会自动考虑偏移吗?

【问题讨论】:

  • kafkadetails 到底是什么?那个物体是什么?来自哪个图书馆?
  • 是KafkaUtils,我会在描述中更正
  • 感谢您的更新,从您所说的我可以看出您使用的是 Spark 2.2 或更低版本,因为 KafkaUtils 在该版本之后不可用。另一个问题是这里的sparkSparkContext 还是SparkSession,因为它们提供了不同的API。
  • 我使用的是 spark 2.4
  • 例如:val stream = KafkaUtils.createDirectStream[String, String](streamingContext, PreferConsistent, Subscribe[String, String](topics, kafkaParams))

标签: scala apache-spark apache-kafka


【解决方案1】:

如果你看看 documentation of kafka connector for Spark,你可以找到大部分答案。

关于 Kafka 连接器的 startingOffsets 选项的文档,最后一部分是关于流式查询。

查询开始时的起点,“earliest”是从最早的偏移量开始,“latest”是从最近的偏移量开始,或者是一个json字符串,为每个TopicPartition指定一个起始偏移量。在json中,-2作为偏移量可以用来表示最早,-1表示最新。注意:对于批量查询,最新(隐式或在 json 中使用 -1)是不允许的。对于流式查询,这仅适用于新查询开始时,并且恢复将始终从查询停止的地方开始。查询期间新发现的分区将最早开始。

如果你有偏移量,它总是会选择可用的偏移量,否则它会向 Kafka 询问 earliestlatest 偏移量。这应该适用于两种类型的流,直接流和结构化流都应该考虑偏移量。

我看到你提到了enable.auto.commit 选项,我只是想确保你知道我上面提供的同一个文档站点中的以下引用。

请注意,以下 Kafka 参数不能设置,Kafka 源或接收器会抛出异常:enable.auto.commit: Kafka 源不提交任何偏移量。

【讨论】:

    【解决方案2】:

    这是我试图回答你的问题

    1. 问题 1:enable.auto.commit 是与 kafka 相关的参数,如果设置为 false,则需要您手动提交(读取更新)您的偏移量到检查点目录。如果您的应用程序重新启动,它将查看检查点目录并从上次提交的偏移量 + 1 开始读取。jaceklaskowski 在这里提到了https://jaceklaskowski.gitbooks.io/apache-kafka/content/kafka-properties-enable-auto-commit.html。作为 spark 应用程序的一部分,无需在任何地方指定偏移量。您所需要的只是检查点目录。此外,请记住,每个消费者的主题中的每个分区都会维护偏移量,因此期望开发人员/用户提供偏移量对 Spark 来说是不好的。
    2. 问题 2:spark.readStream 是一种从 tcp 套接字、kafka 主题等流式源中读取数据的通用方法,而 kafkaUtils 是用于将 spark 与 kafka 集成的专用类,因此我认为如果您使用它会更优化kafka 主题作为来源。我通常自己使用 KafkaUtils,因为我没有做过任何性能基准测试。如果我没记错的话,KafkaUtils 也可以用于订阅多个主题,而 readStream 则不能。
    3. 问题 3:最早的偏移量意味着您的消费者将从可用的最旧记录开始读取,例如,如果您的主题是新的(未发生清理)或未为该主题配置清理,它将从偏移量 0 开始。已配置案例清理,并且删除了直到偏移量 2000 的所有记录,将从偏移量 2001 读取记录,而主题可能有直到偏移量 10000 的记录(这是假设只有一个分区,在主题中将有多个分区,偏移量值为不同的 )。有关更多详细信息,请参阅此处https://spark.apache.org/docs/2.2.0/structured-streaming-kafka-integration.html 的批量查询部分。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-04-27
      • 1970-01-01
      • 2018-06-08
      • 1970-01-01
      • 1970-01-01
      • 2016-11-13
      • 2016-08-21
      • 2019-10-11
      相关资源
      最近更新 更多