【发布时间】: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 guide 在failOnDataLoss 选项下提及:
当数据可能丢失时是否使查询失败(例如, 主题被删除,或偏移量超出范围)。 这可能是假的 警报。当它不能按预期工作时,您可以禁用它。 批处理 如果无法从 由于丢失数据而提供的偏移量。
但我找不到任何关于何时这可以被认为是误报的更多信息,所以我不知道将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.bytes、log.cleaner.enable、log.cleaner.min.compaction.lag.ms和cleanup.policy。我们遇到了类似的问题,调整上述属性给了我们预期的结果failOnDataLosstrue -
我没有设置
log.retention.bytes,所以它应该只根据保留期删除日志。根据cloudurable.com/blog/kafka-architecture-log-compaction/…,日志压缩不应该改变偏移量? -
我遇到了类似的问题,偏移量会突然跳回以前的方式。仍然无法解释这一点。
-
@linehrr 看看我的回答
-
@lfk 是的,我们尝试了这些,但仍在发生。
标签: apache-spark apache-kafka spark-structured-streaming spark-kafka-integration