Spark 2.3.1结构化流状态存储内部工作

Maa*_*mon 8 apache-spark spark-structured-streaming

我一直在浏览关于结构化流的spark 2.3.1的文档,但是找不到有关状态存储在状态存储内部如何工作的详细信息。更具体地说,我想知道的是:(1)状态存储区是分布式的吗?(2)如果是,那么每个工人或每个核心如何?

似乎在旧版本的spark中是每个工人,但现在还不知道。我知道它得到了HDFS的支持,但是没有任何东西可以解释内存中存储的实际工作方式。

确实是分布式内存存储吗?我对重复数据删除特别感兴趣,如果数据来自一个大型数据集,那么这需要进行计划,因为所有“不同”数据集最终都将保存在内存中,直到该数据集处理结束。因此,需要根据状态存储的工作方式来计划工作者或主服务器的大小。

没有人有一些信息,指针或建议如何处理吗?

谢谢,Maatari

cha*_*ash 4

结构化流中只有一种状态存储实现,由内存 HashMap 和 HDFS 支持。In-Memory HashMap 用于数据存储,而 HDFS 用于容错。HashMap占用worker上的执行器内存,每个HashMap代表聚合分区的版本化键值数据(在重复数据删除、groupByy等聚合器运算符之后生成)

但这并不能解释 HDFSBackedStateStore 实际上是如何工作的。我在文档中没有看到它

你是对的,没有这样的文档可用。我必须理解代码(2.3.1),写了一篇关于状态存储在结构化流内部如何工作的文章。您可能想看看:https://www.linkedin.com/pulse/state-management-spark-structed-streaming-chandan-prakash/