【问题标题】:How to implement exactly-once processing when reading reading from directory using Spark Structured Streaming?使用 Spark Structured Streaming 从目录读取内容时如何实现一次性处理?
【发布时间】:2019-02-25 05:48:14
【问题描述】:
我想使用流处理的概念从本地目录读取文件,然后发布到 Apache Kafka。我考虑过使用 Spark Structured Streaming。
在读取 50 行文件后流式传输失败时如何实现检查点。下次启动时是从文件的第 51 行开始,还是从文件的开头再次读取?
另外,如果我们在结构化流中使用检查点,当代码有任何升级或任何更改时,我们会遇到什么问题吗?
【问题讨论】:
标签:
apache-spark
apache-kafka
spark-structured-streaming
【解决方案1】:
在读取 50 行文件后流式传输失败时。下次启动时是从文件的第51行开始,还是从文件的开头再次读取。
要么完全处理整个文件,要么根本不处理。这就是 FileFormat 在 Spark SQL 中的一般工作方式,尤其与 Spark Structured Streaming 几乎没有关系(因为它们共享底层执行基础设施)。
简而言之,引擎“会再次从文件的开头读取。”
也就是说,在 Spark Structured Streaming 中处理文件时没有单行的概念。您一次处理一个作为整个文件(甚至是几个文件)的流式 DataFrame,而您是要逐行处理还是完整处理数据集,这取决于您,Spark 开发人员。
另外,如果我们在结构化流中使用检查点,当代码有任何升级或任何更改时,我们会遇到什么问题吗?
理论上,您不应该这样做。 Spark Structured Streaming 中新的检查点机制(与传统的 Spark Streaming 相比)的目的是允许以更舒适的方式重新启动和升级。检查点仅使用少量信息(通常存储在 JSON 文件中)从最后一个成功的检查点开始重新开始处理。