【问题标题】:Exactly once semantics in Spark Streaming Direct ApproachSpark Streaming Direct Approach 中的 Exactly once 语义
【发布时间】:2020-01-26 15:25:02
【问题描述】:

Spark 的官方文档说,基于 Direct 的方法涉及使用 SimpleConsumer API,它不使用 Zookeeper 来存储偏移量,而是使用 Spark 的元数据检查点来存储偏移量。该文档还说,基于 Direct 的方法保证了语义精确一次。

当我们使用 ssc.checkpoint("directory") 启用 Spark 的元数据检查点时,我们从不指定间隔。

现在,对于每个微批次,在微批次间隔后触发,驱动程序将偏移量发送到每个任务,这些任务为相应的 Kafka 分区检索数据。

问题

  1. 考虑到从 Kafka 检索到的指定偏移量的相应数据不会保留在 Spark 中,并且只有偏移量作为其元数据检查点的一部分存储在 Spark 中,检查点的时间并不重要,因为它直接影响恰好一次或至少/最多一次语义?它是在触发微批处理并且 directstream 从 kafka 检索数据时立即发生,还是在微批处理完成结束时发生?

  2. 另外,作为元数据检查点的一部分存储的偏移量意味着什么?它是否指定已处理的偏移量或尚未处理的偏移量?

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    Checkpointing 是 [CheckpointsKafka itselfYour own data store] 三个选项之一,Checkpointing 有几个缺点,并且不能保证完全一次,除非您的事务是幂等的。

    文档警告您有关检查点的信息如下:

    因此,如果您想要完全一次性语义的等价物,您必须 要么在幂等输出之后存储偏移量,要么将偏移量存储在 与输出并列的原子事务。

    请参阅官方文档的this 部分,详细描述这三个选项

    【讨论】:

    • 没错。我也是这么想的。我的意思是仅启用检查点的基于直接的方法无法在完全一次语义上提供帮助,因为这涉及以原子方式保存处理过的数据和偏移量。而且由于 spark 中的元数据检查点和使用 .saveas*** 保存处理过的数据是两个独立的步骤,因此它不能保证完全一次语义。现在,当您说启用检查点时处理幂等性(忽略检查点的其他问题)可以保证只发生一次语义,您的意思是元数据检查点仅在微批处理完成后发生吗?
    • 我之所以这么说,是因为幂等性来自至少一次语义作为开发人员的责任,这意味着元数据检查点会在每次微批处理执行“之后”发生,以便在微批处理执行后驱动程序立即停机的情况下已完成但在元数据检查点完成之前意味着重新启动驱动程序将意味着重新执行特定的偏移量,因为它上次没有保存在元数据检查点中。我的理解对吗?
    • @SheelPancholi,您实际上可以控制检查点在之前或之后发生,使用参数-“检查点操作符的每个急切标志可以是急切或懒惰的。急切检查点是默认检查点,并立即发生请求时。延迟检查点不会也只会在执行操作时发生“ - jaceklaskowski.gitbooks.io/mastering-spark-sql/…
    猜你喜欢
    • 2019-08-25
    • 2019-12-10
    • 1970-01-01
    • 1970-01-01
    • 2018-01-13
    • 1970-01-01
    • 2019-08-25
    • 2021-12-11
    • 2019-06-23
    相关资源
    最近更新 更多