【问题标题】:Databricks: Queries with streaming sources must be executed with writeStream.start()Databricks:必须使用 writeStream.start() 执行带有流式源的查询
【发布时间】:2021-10-13 02:44:54
【问题描述】:

我知道关于此错误消息还有其他几个问题,但似乎没有一个与我目前面临的问题有关。我正在从 JSON 文件流式传输(这部分有效):

gamingEventDF = (spark
.readStream
.schema(eventSchema) 
.option('streamName','mobilestreaming_demo') 
.option("maxFilesPerTrigger", 1)
.json(inputPath) 
)

接下来我想使用 writeStream 将其附加到表中:

def writeToBronze(sourceDataframe, bronzePath, streamName):
(sourceDataframe.rdd
.spark
.writeStream.format("delta")
.option("checkpointLocation", bronzePath + "/_checkpoint")
.queryName(streamName)
.outputMode("append") 
.start(bronzePath)
)

当我现在跑步时:

writeToBronze(gamingEventDF, outputPathBronze, "bronze_stream")

我收到错误:AnalysisException: 必须使用 writeStream.start() 执行带有流源的查询

顺便说一句:当我删除 .rdd 时,我收到另一个错误('DataFrame' 对象没有属性'spark')

知道我做错了什么吗? 非常感谢

【问题讨论】:

    标签: apache-spark pyspark spark-streaming


    【解决方案1】:

    writeStream 方法在数据帧类上可用,而不是在 SparkSession 上。

    下面的代码应该适合你。

    def writeToBronze(sourceDataframe, bronzePath, streamName):
      (sourceDataframe
      .writeStream.format("delta")
      .option("checkpointLocation", bronzePath + "/_checkpoint")
      .queryName(streamName)
      .outputMode("append") 
      .start(bronzePath)
      .awaitTermination())
    

    【讨论】:

      猜你喜欢
      • 2017-03-29
      • 1970-01-01
      • 2021-01-31
      • 2018-03-14
      • 1970-01-01
      • 2021-08-16
      • 1970-01-01
      • 2019-05-25
      • 2017-06-23
      相关资源
      最近更新 更多