【发布时间】:2021-11-22 05:47:18
【问题描述】:
我有一个 Spark 流处理器。 数据框 dfNewExceptions 有重复项(由“ExceptionId”重复)。 由于这是一个流式数据集,因此以下查询失败:
val dfNewUniqueExceptions = dfNewExceptions.sort(desc("LastUpdateTime"))
.coalesce(1)
.dropDuplicates("ExceptionId")
val dfNewExceptionCore = dfNewUniqueExceptions.select("ExceptionId", "LastUpdateTime")
dfNewExceptionCore.writeStream
.format("console")
// .outputMode("complete")
.option("truncate", "false")
.option("numRows",5000)
.start()
.awaitTermination(1000)
** 线程“主”org.apache.spark.sql.AnalysisException 中的异常:流数据帧/数据集不支持排序,除非它在完整输出模式下聚合的数据帧/数据集上;; **
关于如何从 dfNewExceptions 中删除重复项有什么建议吗?
【问题讨论】:
标签: scala apache-spark apache-spark-sql spark-structured-streaming delta-lake