【问题标题】:How to reset spark structured streaming data to last available offset如何将 Spark 结构化流数据重置为最后可用的偏移量
【发布时间】:2020-08-01 04:24:48
【问题描述】:

我正在使用 Kafka 运行结构化流应用程序。我发现如果由于某种原因系统关闭了几天......检查点变得陈旧并且在Kafka中找不到与检查点相对应的偏移量。如何让 Spark Structured Streaming 应用程序选择最后一个可用的偏移量并从那里开始。我尝试将偏移重置设置为较早/最新,但系统因以下错误而崩溃:

org.apache.kafka.clients.consumer.OffsetOutOfRangeException: Offsets out of range with no configured reset policy for partitions: {MyTopic-574=6559828}
at org.apache.kafka.clients.consumer.internals.Fetcher.parseCompletedFetch(Fetcher.java:970)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:490)
at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1259)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1187)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1115)
at org.apache.spark.sql.kafka010.InternalKafkaConsumer.fetchData(KafkaDataConsumer.scala:470)
at org.apache.spark.sql.kafka010.InternalKafkaConsumer.org$apache$spark$sql$kafka010$InternalKafkaConsumer$$fetchRecord(KafkaDataConsumer.scala:361)
at org.apache.spark.sql.kafka010.InternalKafkaConsumer$$anonfun$get$1.apply(KafkaDataConsumer.scala:251)
at org.apache.spark.sql.kafka010.InternalKafkaConsumer$$anonfun$get$1.apply(KafkaDataConsumer.scala:234)
at org.apache.spark.util.UninterruptibleThread.runUninterruptibly(UninterruptibleThread.scala:77)
at org.apache.spark.sql.kafka010.InternalKafkaConsumer.runUninterruptiblyIfPossible(KafkaDataConsumer.scala:209)
at org.apache.spark.sql.kafka010.InternalKafkaConsumer.get(KafkaDataConsumer.scala:234)
at org.apache.spark.sql.kafka010.KafkaDataConsumer$class.get(KafkaDataConsumer.scala:64)
at org.apache.spark.sql.kafka010.KafkaDataConsumer$CachedKafkaDataConsumer.get(KafkaDataConsumer.scala:500)
at org.apache.spark.sql.kafka010.KafkaMicroBatchInputPartitionReader.next(KafkaMicroBatchReader.scala:357)
at org.apache.spark.sql.execution.datasources.v2.DataSourceRDD$$anon$1.hasNext(DataSourceRDD.scala:49)
at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:409)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.processNext(Unknown Source)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
at org.apache.spark.sql.execution.UnsafeExternalRowSorter.sort(UnsafeExternalRowSorter.java:216)
at org.apache.spark.sql.execution.SortExec$$anonfun$1.apply(SortExec.scala:108)
at org.apache.spark.sql.execution.SortExec$$anonfun$1.apply(SortExec.scala:101)
at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:836)
at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:836)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
at org.apache.spark.scheduler.Task.run(Task.scala:123)
at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)

【问题讨论】:

  • 你能把你的消费者的代码和你的配置一起显示吗?另外,您使用的是哪个版本?
  • 我遇到了同样的问题,我遵循了这种方法 - 正如您提到的检查点中的偏移量太旧,请您备份并删除您的检查点位置并尝试启动您的 spark 应用程序。那应该工作..
  • @Srinivas 非常感谢您的回复。我们不想触及实时流管道......如果偏移量是旧的,我们希望 Spark 忽略偏移量并继续。我已将 failOnDataLoss 设置为 false,但这并没有解决我的问题。如果发生上述问题,我必须清理检查点,即用户数据,然后重新启动 Spark。这对于生产系统是不可接受的。

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


【解决方案1】:

如果系统宕机了几天,则意味着某些日志可能已被压缩。更准确地说,您的应用程序尝试从主题MyTopic 中的第574 个分区读取偏移量6559828


为了找到每个分区的最早可用偏移量,您可以简单地运行以下命令:

bin/kafka-run-class.sh kafka.tools.GetOffsetShell \
    --broker-list localhost:9092 \
    --topic MyTopic \
    --time -2

【讨论】:

  • 非常感谢您的评论。我希望结构化流应用程序不会停止...继续处理而不使用它可以找到的最后一个偏移量。
  • @Alchemist 你只需要auto.offset.resetlatest。应用程序损坏的原因是发生了日志压缩并删除了一些消息,并且您的流应用程序找不到从哪里开始。还可以考虑增加日志压缩周期,以避免将来出现这种行为。
  • 谢谢@Giorgos Myrianthous...一个月来我一直在努力调整这个简单的结构化流应用程序...我已经将 auto.offset.reset 设置为最新的。在调整使用 Kafka 的 spark 系统时,您是否有任何我需要绝对考虑的调整参数列表...我正在使用保留设置为 3 天的 Kafka 主题。我无法修改主题。看起来这是影响该主题的其他消费者的主题级别设置。但是,对于使用 Kafka 的 Spark Structured Streaming,是否有任何推荐的设置...非常感谢您的回复。
  • @Alchemist 我没有建议您修改主题并影响其他消费者。我说的是你的流应用程序。
  • 谢谢我已经使用startingOffsets将auto.offset.reset设置为最新的spark。但不知道如何在 Spark Structured Streaming 消费者上设置日志压缩周期。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-09-28
  • 2018-11-23
  • 2021-03-18
  • 2019-10-03
  • 1970-01-01
  • 2019-12-25
相关资源
最近更新 更多