Spark3.2.2(Scala):过滤DataFrame中null值并记录日志
解决方案
步骤说明
- 先提取所有
weight为null的品牌,收集到Driver端后打印指定格式的提示信息 - 过滤原DataFrame,剔除
weight为null的无效记录
完整代码示例
import org.apache.spark.sql.SQLContext import org.apache.spark.sql.functions.col // 假设你的SQLContext已初始化,原DataFrame名为df // 1. 获取weight为null的品牌并打印日志 val nullWeightBrands = df.filter(col("weight").isNull) .select("brand") .as[String] .collect() nullWeightBrands.foreach(brand => { println(s"Weight for brand $brand is null, dropping it from the data") }) // 2. 生成过滤后的有效DataFrame val filteredDf = df.filter(col("weight").isNotNull) // 可选:验证过滤结果 filteredDf.show()
代码细节说明
col("weight").isNull:精准匹配weight字段为null的记录.as[String].collect():将品牌字段转为字符串类型并收集到Driver端,确保日志在Driver统一输出(避免分布式场景下Executor重复打印)- Scala字符串插值
s"...":快速生成符合要求的提示消息格式 col("weight").isNotNull:过滤保留weight不为null的有效数据
大数据量场景适配
如果原DataFrame数据量极大,collect()可能导致Driver内存溢出,可改用foreachPartition在Executor端分散打印日志:
df.filter(col("weight").isNull) .select("brand") .foreachPartition(iter => { iter.foreach(row => { val brand = row.getString(0) println(s"Weight for brand $brand is null, dropping it from the data") }) })
内容的提问来源于stack exchange,提问作者Nab
相关产品推荐
相关产品推荐

