【发布时间】: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