Spark Streaming DataFrame按ExceptionId去重报错解决方案咨询
你遇到的报错是因为Structured Streaming不支持在无界流数据集上直接执行全局排序操作,仅聚合后的DataFrame在Complete输出模式下允许排序。你原有逻辑中排序是为了去重时保留同一个ExceptionId对应的最新LastUpdateTime记录,可通过以下两种方案实现需求:
方案1:分组聚合取最新记录(匹配保留最新记录的需求)
直接按ExceptionId分组取LastUpdateTime最大值,天然满足去重且保留最新记录的要求,适配流处理场景:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.streaming.OutputMode // 可选:添加水位线控制状态存储时长,避免状态无限膨胀,时长可根据业务延迟情况调整 val dfWithWatermark = dfNewExceptions .withWatermark("LastUpdateTime", "24 hours") val dfNewUniqueExceptions = dfWithWatermark .groupBy("ExceptionId") .agg(max("LastUpdateTime").alias("LastUpdateTime")) val dfNewExceptionCore = dfNewUniqueExceptions.select("ExceptionId", "LastUpdateTime") dfNewExceptionCore.writeStream .format("console") .outputMode(OutputMode.Append()) // 已添加水位线用Append模式,未加水印可替换为Update模式 .option("truncate", "false") .option("numRows",5000) .start() .awaitTermination(1000)
方案2:使用流原生dropDuplicates算子(适合仅需去重、不需要保留最新记录的场景)
Structured Streaming原生支持dropDuplicates算子,会自动维护已出现过的ExceptionId状态,默认保留第一条出现的记录:
import org.apache.spark.sql.functions._ // 同样可添加水位线控制状态大小 val dfWithWatermark = dfNewExceptions .withWatermark("LastUpdateTime", "24 hours") .dropDuplicates("ExceptionId") // 后续输出逻辑和方案1一致
注意事项
不要在流处理逻辑中使用coalesce(1),该操作会强制把所有数据合并到单分区处理,会严重降低流处理性能,甚至引发作业阻塞。
内容的提问来源于stack exchange,提问作者Ganesha
相关产品推荐
相关产品推荐

