【问题标题】:Spark Structured Streaming Kafka error -- offset was changedSpark Structured Streaming Kafka 错误——偏移量已更改
【发布时间】:2019-07-30 00:48:23
【问题描述】:

我的 Spark Structured Streaming 应用程序运行了几个小时,然后出现此错误而失败

java.lang.IllegalStateException: Partition [partition-name] offset was changed from 361037 to 355053, some data may have been missed.
Some data may have been lost because they are not available in Kafka any more; either the
data was aged out by Kafka or the topic may have been deleted before all the data in the
topic was processed. If you don't want your streaming query to fail on such cases, set the
source option "failOnDataLoss" to "false".

偏移量当然每次都不同,但第一个总是大于第二个。主题数据不能过期,因为主题的保留期是 5 天,我昨天重新创建了这个主题,但今天又出现错误。从中恢复的唯一方法是删除检查点。

Spark's Kafka integration guidefailOnDataLoss 选项下提及:

当数据可能丢失时是否使查询失败(例如, 主题被删除,或偏移量超出范围)。 这可能是假的 警报。当它不能按预期工作时,您可以禁用它。 批处理 如果无法从 由于丢失数据而提供的偏移量。

但我找不到任何关于何时这可以被认为是误报的更多信息,所以我不知道将failOnDataLoss 设置为false 是否安全,或者是否存在我的集群存在实际问题(在这种情况下,我们实际上会丢失数据)。

更新:我调查了 Kafka 日志,在 Spark 失败的所有情况下,Kafka 都记录了几条这样的消息(我假设每个 Spark 消费者都有一条消息):

INFO [GroupCoordinator 1]: Preparing to rebalance group spark-kafka-...-driver-0 with old generation 1 (__consumer_offsets-25) (kafka.coordinator.group.GroupCoordinator)

【问题讨论】:

  • 有保留期,你应该检查kafka的log.retention.byteslog.cleaner.enablelog.cleaner.min.compaction.lag.mscleanup.policy。我们遇到了类似的问题,调整上述属性给了我们预期的结果failOnDataLosstrue
  • 我没有设置log.retention.bytes,所以它应该只根据保留期删除日志。根据cloudurable.com/blog/kafka-architecture-log-compaction/…,日志压缩不应该改变偏移量?
  • 我遇到了类似的问题,偏移量会突然跳回以前的方式。仍然无法解释这一点。
  • @linehrr 看看我的回答
  • @lfk 是的,我们尝试了这些,但仍在发生。

标签: apache-spark apache-kafka spark-structured-streaming spark-kafka-integration


【解决方案1】:

我没有这个问题了。我做了两处改动:

  1. 禁用 YARN 的动态资源分配(这意味着我必须手动计算最佳执行器数量等并将它们传递给spark-submit
  2. 升级到 Spark 2.4.0,同时将 Kafka 客户端从 0.10.0.1 升级到 2.0.0

禁用动态资源分配意味着执行程序(=消费者)不会在应用程序运行时创建和终止,从而无需重新平衡。所以这很可能是阻止错误发生的原因。

【讨论】:

    【解决方案2】:

    这似乎是旧版本 Spark 和 spark-sql-kafka 库中的一个已知错误。

    我发现以下 JIRA 票证相关:

    • SPARK-28641: MicroBatchExecution 提交的偏移量大于可用的偏移量
    • SPARK-26267: Kafka 源可能会重新处理数据
    • KAFKA-7703: KafkaConsumer.position 可能在调用“seekToEnd”后返回错误的偏移量

    简而言之,引用参与其中的开发人员:

    “这是 Kafka 中的一个已知问题,请参阅 KAFKA-7703。这已在 SPARK-26267 中的 2.4.1 和 3.0.0 中修复。请将 Spark 升级到更高版本。另一种可能性是将 Kafka 升级到 2.3 .0,其中 Kafka 端是固定的。”

    “KAFKA-7703 仅存在于 Kafka 1.1.0 及更高版本中,因此可能的解决方法是使用没有此问题的旧版本。这不会影响 Spark 2.3.x 及更低版本,因为我们使用的是 Kafka 0.10默认为 .0.1。”

    在我们的案例中,我们在 HDP 3.1 平台上遇到了同样的问题。我们有 Spark 2.3.2 和 spark-sql-kafka 库 (https://mvnrepository.com/artifact/org.apache.spark/spark-sql-kafka-0-10_2.11/2.3.2.3.1.0.0-78),但是,使用 kafka-clients 2.0.0。这意味着我们由于后续条件而面临这个错误:

    • 我们的火花
    • 1.1.0

    变通解决方案

    我们能够通过删除包含0 偏移量的批次号的“偏移量”子文件夹中的检查点文件来解决此问题。

    删除此文件时,请确保子文件夹“commits”和“offset”中的检查点文件中的批号在删除后仍然匹配。

    【讨论】:

      猜你喜欢
      • 2021-05-22
      • 2020-09-03
      • 2018-04-26
      • 2017-12-31
      • 1970-01-01
      • 1970-01-01
      • 2017-06-22
      • 2017-02-06
      • 2018-09-22
      相关资源
      最近更新 更多