Chi*_*bon 3 scala apache-spark spark-structured-streaming
运行 spark 作业时出现以下错误:
org.apache.spark.sql.AnalysisException:当流数据帧/数据集上有流聚合时,不支持追加输出模式;;
我不确定问题是否是由于缺少水印引起的,我不知道如何在这种情况下应用。以下是应用的聚合操作:
def aggregateByValue(): DataFrame = {
df.withColumn("Value", expr("(BookingClass, Value)"))
.groupBy("AirlineCode", "Origin", "Destination", "PoS", "TravelDate", "StartSaleDate", "EndSaleDate", "avsFlag")
.agg(collect_list("Value").as("ValueSeq"))
.drop("Value")
}
Run Code Online (Sandbox Code Playgroud)
用法:
val theGroupedDF = theDF
.multiplyYieldByHundred
.explodeDates
.aggregateByValue
val query = theGroupedDF.writeStream
.outputMode("append")
.format("console")
.start()
query.awaitTermination()
Run Code Online (Sandbox Code Playgroud)
更改outputMode以complete解决问题。
val query = theGroupedDF.writeStream
.outputMode("complete")
.format("console")
.start()
query.awaitTermination()
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
2874 次 |
| 最近记录: |