【问题标题】:How to ensure no data loss for kafka data ingestion through Spark Structured Streaming?如何通过 Spark Structured Streaming 确保 kafka 数据摄取不丢失数据?
【发布时间】:2020-12-24 08:30:49
【问题描述】:

我有一个长期运行的 Spark 结构化流式传输作业,它正在摄取 kafka 数据。我有一个担忧如下。如果作业由于某种原因失败并稍后重新启动,如何确保 kafka 数据将从断点开始摄取,而不是在作业重新启动时始终摄取当前和以后的数据。我是否需要明确指定消费者组和 auto.offet.reset 等内容?他们是否支持 spark kafka 摄取?谢谢!

【问题讨论】:

  • 谢谢,我说的是消费。只想设置消费者组 id 以确保在 spark 作业失败时保留偏移量。我的火花是2.4.6。 kafka lib是0.10。当我设置 group.id 时,我收到以下错误“线程“主”java.lang.IllegalArgumentException 中的异常:不支持 Kafka 选项“group.id”,因为用户指定的消费者组不用于跟踪偏移量。”

标签: apache-spark apache-kafka spark-streaming kafka-consumer-api


【解决方案1】:

根据Spark Structured Integration Guide,Spark 本身会跟踪偏移量,并且没有将偏移量提交回 Kafka。这意味着如果您的 Spark Streaming 作业失败并且您重新启动它,所有关于偏移量的必要信息都存储在 Spark 的检查点文件中。这样您的应用程序就会知道它在哪里停止并继续处理剩余的数据。

我已经在另一个post 中写了更多关于设置group.id 和Spark 的偏移检查点的详细信息

以下是针对 Spark 结构化流作业的最重要的 Kafka 特定配置:

group.id:Kafka 源将为每个查询自动创建一个唯一的组 ID。根据代码,group.id 将自动设置为

val uniqueGroupId = s"spark-kafka-source-${UUID.randomUUID}-${metadataPath.hashCode}

auto.offset.reset:设置源选项 startingOffsets 以指定从哪里开始。 Structured Streaming 管理内部消耗的偏移量,而不是依赖 kafka Consumer 来完成

enable.auto.commit:Kafka 源不提交任何偏移量。

因此,在结构化流中,目前无法为 Kafka 消费者定义您的自定义 group.id,并且结构化流在内部管理偏移量,而不是提交回 Kafka(也不是自动)。

【讨论】:

  • 谢谢迈克,但我还是有些困惑。看链接spark.apache.org/docs/latest/…,里面明确提到'kafka.group.id'是可以设置的,只是需要非常小心。我想知道 Kafka 消费者的自定义 group.id 是否不可能,或者只是在某些最新版本中可能,例如 3.0.0。谢谢The Kafka group id to use in Kafka consumer while reading from Kafka. Use this with caution. By default, each query generates a unique group id for reading data
  • 是的,仅适用于 v3,不适用于 2.4.6。另请查看您在其他question 中针对类似主题提出的答案。
  • 谢谢迈克。是的,我之前也问过一个类似的话题,但还不是很确定。我是否可以理解,如果以下假设对于在 v3 中设置 kafka.group.id 是正确的? 1)Kafka经纪人将能够自己维护偏移量作为kafka标准。 2)自定义的group_id会覆盖spark维护的内部group_id提交给kafka broker。 3) Spark 将自动提交回 kafka。我不需要做任何其他事情来避免数据丢失,例如手动提交偏移量等,对吧?
  • 感谢您的提醒。我是 Stackoverflow 的新手,所以我忘了接受答案。是真的。根据您的 cmets,我现在已经接受了它们。对于大多数问题,我之前也投过赞成票。对我来说,我确实尊重人们在回答我的问题时的帮助。至于我下面给你的问题,我也是从“谢谢”开始的。至于类似的问题,我也解释说只是为了进一步确认。即使作为您的回答,您也提到“目前无法定义您的自定义 group.id”,所以我想再次确认 spark 3.0 是否支持它,以及它是如何支持的。
猜你喜欢
  • 1970-01-01
  • 2019-11-20
  • 2017-04-25
  • 1970-01-01
  • 1970-01-01
  • 2021-05-22
  • 2020-07-25
  • 2019-06-28
  • 2018-09-06
相关资源
最近更新 更多