结构化流异常:流聚合不支持追加输出模式

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)

Chi*_*bon 7

更改outputModecomplete解决问题。

val query = theGroupedDF.writeStream
  .outputMode("complete")
  .format("console")
  .start()
query.awaitTermination()
Run Code Online (Sandbox Code Playgroud)