【问题标题】:How to achieve exactly-once write guaranty with foreachBatch sink in Spark Structured Streaming如何在 Spark Structured Streaming 中使用 foreachBatch sink 实现一次性写入保证
【发布时间】:2020-10-18 04:28:03
【问题描述】:

来自docs

默认情况下,foreachBatch 仅提供至少一次写入保证。 但是,您可以使用提供给函数的 batchId 作为 对输出进行重复数据删除并获得完全一次的保证。

这是否意味着除非我们加倍努力,否则即使我们在主 writeStream 操作中使用检查点,我们也无法实现一次性写入保证

如果是,应该怎么做才能实现一次性写入保证?文档中的含义:

使用批处理 ID ?

PS:最初的问题是针对 Kafka 的,但我将其概括为解决方案适用于 foreachBatch 块中的任何内容。

【问题讨论】:

    标签: apache-spark apache-kafka apache-spark-sql spark-structured-streaming exactly-once


    【解决方案1】:

    这是否意味着除非我们加倍努力,否则即使我们在主 writeStream 操作中使用检查点,我们也无法实现一次性写入保证?

    结构化流保证至少有一个语义,这意味着每条记录将至少出现一次,即使我们有检查点,也不能保证不会出现重复记录。

    如果是,应该怎么做才能实现一次性写入保证?文档中的含义是什么

    根据选择使用的数据接收器,实现恰好一次语义的方式会有所不同。

    为了便于解释,我们将弹性搜索作为数据接收器。 我们知道 ES 是一个文档存储,每条记录都有一个唯一的 doc_id。 假设我们有一个具有以下架构的数据框 -

    |-- count: long (nullable = false)
    |-- department: string (nullable = false)
    |-- start_time: string (nullable = false)
    |-- end_time: string (nullable = false)
    

    在这种情况下,我们将 (department,start_time,end_time) 作为唯一的键列,我们可以找到这些列的哈希并将其用作弹性搜索中的 doc_id 列,使用index-api

    使用这种方式,在重复记录的情况下,它们将被散列到相同的值,并且由于 ES 不允许在特定索引中重复 doc_id,它将使用相同的记录更新文档并增加其版本。

    其他数据接收器也可以采用类似的方法。

    【讨论】:

      猜你喜欢
      • 2018-01-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-04-04
      • 2023-03-13
      • 2020-05-05
      相关资源
      最近更新 更多