【问题标题】: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,它将使用相同的记录更新文档并增加其版本。
其他数据接收器也可以采用类似的方法。