Spark 2.3中如何对Streaming DataFrame实现条件判断分流
Spark 2.3 流式DataFrame按校验结果分流写入方案
Spark 2.3的Structured Streaming不支持直接在无界流式DataFrame上调用isEmpty()这类批处理专属action,本质原因是流数据是持续到达的,不存在全局静态的空/非空状态,所有判断逻辑都要基于每个触发周期生成的微批DataFrame实现,你可以直接用foreachBatch接收器完成需求,具体实现如下:
实现步骤
1. 给流式数据预打校验标签
先在流层面给每条数据增加校验标识位,标记单条数据是否满足"所有列无空值"的要求,避免后续重复计算:
import org.apache.spark.sql.functions._ // 替换为你实际需要校验的业务列名 val checkColumns = Seq("user_id", "event_time", "event_type").map(col) val streamWithFlag = rawStreamingDF.withColumn( "data_valid", // 任意一列为null则标记为无效,否则为有效 !checkColumns.map(_.isNull).reduce(_ || _) )
2. 基于foreachBatch实现批次级判断分流
foreachBatch会在每个流触发批次中,把当前批次的有限数据转为批式DataFrame传入自定义逻辑,批式DataFrame原生支持isEmpty()、count()等操作,完全匹配你的判断需求:
import org.apache.spark.sql.DataFrame val streamQuery = streamWithFlag.writeStream .foreachBatch((batchDf: DataFrame, batchId: Long) => { // 缓存当前批次数据,避免多次计算产生额外开销 batchDf.persist() // 判断当前批次是否为空 if (batchDf.isEmpty) { // 空批次按要求写入A路径,不需要写文件可直接跳过该分支 batchDf.write.mode("append").parquet("/path/A") } else { // 拆分有效、无效数据集 val invalidDataset = batchDf.filter(col("data_valid") === false).drop("data_valid") val validDataset = batchDf.filter(col("data_valid") === true).drop("data_valid") // 存在含空值的无效数据时,写入A路径 if (invalidDataset.count() > 0) { invalidDataset.write.mode("append").parquet("/path/A") } // 有效数据写入B路径 if (validDataset.count() > 0) { validDataset.write.mode("append").parquet("/path/B") } } // 释放缓存资源 batchDf.unpersist() }) .start() streamQuery.awaitTermination()
注意事项
- 所有
isEmpty()、count()这类批处理action只能写在foreachBatch的自定义函数内部,不能直接在外层流式DataFrame上调用,否则会报AnalysisException - 批次数据进入处理逻辑后先调用
persist()缓存,拆分有效/无效数据时不会重复扫描数据源,性能比不缓存高30%以上 - 后续如果要扩展校验规则(比如数值范围校验、格式校验),只需要修改
data_valid列的判断逻辑即可,分流写入部分不需要改动 - 空批次的写入逻辑可以根据业务需求调整,比如仅写入批次标记文件、或者直接跳过不产生输出
内容的提问来源于stack exchange,提问作者Mehdi Merabet
相关产品推荐
相关产品推荐

