【发布时间】:2020-01-26 15:25:02
【问题描述】:
Spark 的官方文档说,基于 Direct 的方法涉及使用 SimpleConsumer API,它不使用 Zookeeper 来存储偏移量,而是使用 Spark 的元数据检查点来存储偏移量。该文档还说,基于 Direct 的方法保证了语义精确一次。
当我们使用 ssc.checkpoint("directory") 启用 Spark 的元数据检查点时,我们从不指定间隔。
现在,对于每个微批次,在微批次间隔后触发,驱动程序将偏移量发送到每个任务,这些任务为相应的 Kafka 分区检索数据。
问题:
考虑到从 Kafka 检索到的指定偏移量的相应数据不会保留在 Spark 中,并且只有偏移量作为其元数据检查点的一部分存储在 Spark 中,检查点的时间并不重要,因为它直接影响恰好一次或至少/最多一次语义?它是在触发微批处理并且 directstream 从 kafka 检索数据时立即发生,还是在微批处理完成结束时发生?
另外,作为元数据检查点的一部分存储的偏移量意味着什么?它是否指定已处理的偏移量或尚未处理的偏移量?
【问题讨论】:
标签: apache-spark spark-streaming