You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark3.2.2(Scala):过滤DataFrame中null值并记录日志

解决方案

步骤说明

  1. 先提取所有weight为null的品牌,收集到Driver端后打印指定格式的提示信息
  2. 过滤原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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 04:23:11