【问题标题】:Why does using cache on streaming Datasets fail with "AnalysisException: Queries with streaming sources must be executed with writeStream.start()"?为什么在流数据集上使用缓存失败并出现“AnalysisException:必须使用 writeStream.start() 执行带有流源的查询”?
【发布时间】:2017-06-23 01:24:37
【问题描述】:
SparkSession
  .builder
  .master("local[*]")
  .config("spark.sql.warehouse.dir", "C:/tmp/spark")
  .config("spark.sql.streaming.checkpointLocation", "C:/tmp/spark/spark-checkpoint")
  .appName("my-test")
  .getOrCreate
  .readStream
  .schema(schema)
  .json("src/test/data")
  .cache
  .writeStream
  .start
  .awaitTermination

在 Spark 2.1.0 中执行此示例时出现错误。 如果没有 .cache 选项,它会按预期工作,但使用 .cache 选项我得到了:

Exception in thread "main" org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start();;
FileSource[src/test/data]
at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.org$apache$spark$sql$catalyst$analysis$UnsupportedOperationChecker$$throwError(UnsupportedOperationChecker.scala:196)
at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$$anonfun$checkForBatch$1.apply(UnsupportedOperationChecker.scala:35)
at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$$anonfun$checkForBatch$1.apply(UnsupportedOperationChecker.scala:33)
at org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala:128)
at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.checkForBatch(UnsupportedOperationChecker.scala:33)
at org.apache.spark.sql.execution.QueryExecution.assertSupported(QueryExecution.scala:58)
at org.apache.spark.sql.execution.QueryExecution.withCachedData$lzycompute(QueryExecution.scala:69)
at org.apache.spark.sql.execution.QueryExecution.withCachedData(QueryExecution.scala:67)
at org.apache.spark.sql.execution.QueryExecution.optimizedPlan$lzycompute(QueryExecution.scala:73)
at org.apache.spark.sql.execution.QueryExecution.optimizedPlan(QueryExecution.scala:73)
at org.apache.spark.sql.execution.QueryExecution.sparkPlan$lzycompute(QueryExecution.scala:79)
at org.apache.spark.sql.execution.QueryExecution.sparkPlan(QueryExecution.scala:75)
at org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:84)
at org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:84)
at org.apache.spark.sql.execution.CacheManager$$anonfun$cacheQuery$1.apply(CacheManager.scala:102)
at org.apache.spark.sql.execution.CacheManager.writeLock(CacheManager.scala:65)
at org.apache.spark.sql.execution.CacheManager.cacheQuery(CacheManager.scala:89)
at org.apache.spark.sql.Dataset.persist(Dataset.scala:2479)
at org.apache.spark.sql.Dataset.cache(Dataset.scala:2489)
at org.me.App$.main(App.scala:23)
at org.me.App.main(App.scala)

有什么想法吗?

【问题讨论】:

  • 抱歉,我不认为不使用缓存就是解决方案。
  • Martin,请随时参加 SPARK-20927 上的 cmets,了解流计算缓存的需求

标签: scala apache-spark apache-spark-sql apache-spark-2.0 spark-structured-streaming


【解决方案1】:

您的(非常有趣的)案例归结为以下行(您可以在spark-shell 中执行):

scala> :type spark
org.apache.spark.sql.SparkSession

scala> spark.readStream.text("files").cache
org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start();;
FileSource[files]
  at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.org$apache$spark$sql$catalyst$analysis$UnsupportedOperationChecker$$throwError(UnsupportedOperationChecker.scala:297)
  at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$$anonfun$checkForBatch$1.apply(UnsupportedOperationChecker.scala:36)
  at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$$anonfun$checkForBatch$1.apply(UnsupportedOperationChecker.scala:34)
  at org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala:127)
  at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.checkForBatch(UnsupportedOperationChecker.scala:34)
  at org.apache.spark.sql.execution.QueryExecution.assertSupported(QueryExecution.scala:63)
  at org.apache.spark.sql.execution.QueryExecution.withCachedData$lzycompute(QueryExecution.scala:74)
  at org.apache.spark.sql.execution.QueryExecution.withCachedData(QueryExecution.scala:72)
  at org.apache.spark.sql.execution.QueryExecution.optimizedPlan$lzycompute(QueryExecution.scala:78)
  at org.apache.spark.sql.execution.QueryExecution.optimizedPlan(QueryExecution.scala:78)
  at org.apache.spark.sql.execution.QueryExecution.sparkPlan$lzycompute(QueryExecution.scala:84)
  at org.apache.spark.sql.execution.QueryExecution.sparkPlan(QueryExecution.scala:80)
  at org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:89)
  at org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:89)
  at org.apache.spark.sql.execution.CacheManager$$anonfun$cacheQuery$1.apply(CacheManager.scala:104)
  at org.apache.spark.sql.execution.CacheManager.writeLock(CacheManager.scala:68)
  at org.apache.spark.sql.execution.CacheManager.cacheQuery(CacheManager.scala:92)
  at org.apache.spark.sql.Dataset.persist(Dataset.scala:2603)
  at org.apache.spark.sql.Dataset.cache(Dataset.scala:2613)
  ... 48 elided

这个原因很容易解释(对 Spark SQL 的explain 没有双关语)。

spark.readStream.text("files") 创建一个所谓的流数据集

scala> val files = spark.readStream.text("files")
files: org.apache.spark.sql.DataFrame = [value: string]

scala> files.isStreaming
res2: Boolean = true

流数据集是 Spark SQL 的Structured Streaming 的基础。

您可能已经阅读了结构化流的Quick Example

然后使用start() 开始流式计算。

引用DataStreamWriter的start的scaladoc:

start(): StreamingQuery 开始执行流式查询,当新数据到达时,它将不断将结果输出到给定路径。

因此,您必须使用start(或foreach)来开始执行流式查询。你已经知道了。

但是……结构化流中有Unsupported Operations

此外,还有一些 Dataset 方法不适用于流数据集。它们是会立即运行查询并返回结果的操作,这在流式数据集上没有意义。

如果您尝试这些操作中的任何一个,您将看到一个 AnalysisException,例如“流数据帧/数据集不支持操作 XYZ”。

这看起来很熟悉,不是吗?

cache不在不支持的操作列表中,但那是因为它被简单地忽略了(我报告了SPARK-20927 来修复它)。

cache 应该在列表中,因为它确实在查询注册到 Spark SQL 的 CacheManager 之前执行查询。

让我们深入了解 Spark SQL...屏住呼吸...

cacheispersistpersistrequests the current CacheManager to cache the query

sparkSession.sharedState.cacheManager.cacheQuery(this)

缓存查询时CacheManager 确实 execute it:

sparkSession.sessionState.executePlan(planToCache).executedPlan

我们知道是不允许的,因为这是start(或foreach)这样做。

问题解决了!

【讨论】:

  • 我认为这是一个错误,所以我更早地报告了它issues.apache.org/jira/browse/SPARK-20865,我只需要确认我的强硬。谢谢。
  • 与 master 的链接并不真正相关,因为目标代码可以更改。我认为这是您的链接中附加的内容
  • 我也想知道,但是由于您的链接在代码更改时没有价值,我建议您针对特定的提交。您的帖子是在 T 时间写的,所以可能与未来的 spark 版本无关。我并不认为您的帖子仅在特定日期是真实的。
  • @mathieu 好的。你说得对。它可能不是由设计支持的。请注意,Spark 2.3 附带另一个流引擎(这可能会改变重新缓存的内容)。
  • @jackek - 你是否仍然对如何链接到 github 中的 master 以外的东西感兴趣?如果你是:在这里。 1) 查看链接到文件“主”版本的页面顶部。您将看到完整路径。右边是 2 个按钮,分别是“查找文件”和“复制路径”。在路径的左侧(这是一系列指向父目录的可点击链接),您将看到一个带有文本“master”的按钮。现在只需单击它,您将看到一个下拉列表,可让您选择要链接到的特定标签或分支。只需选择其中一个以获得较少“不稳定”的链接。
猜你喜欢
  • 1970-01-01
  • 2017-03-29
  • 2018-03-14
  • 2021-10-13
  • 2021-01-31
  • 2021-08-16
  • 1970-01-01
  • 1970-01-01
  • 2020-09-01
相关资源
最近更新 更多