【问题标题】:Spark: replace null values in dataframe with mean of columnSpark:用列的平均值替换数据框中的空值
【发布时间】:2016-11-16 07:47:15
【问题描述】:

如何创建一个 UDF 以编程方式将火花数据框中每列中的空值替换为列平均值。例如,在示例数据 col1 中,空值的值为 ((2+4+6+8+5)/5) = 5。

示例数据:

col1    col2    col3
2       null    3
4       3       3
6       5       null
8       null    2
null    6       4
5       2       8

所需数据:

col1    col2    col3
2       4       3
4       3       3
6       5       4
8       4       2
5       6       4
5       2       8

【问题讨论】:

  • 在纯 SQl 中,这可以通过交叉连接每个列的表并使用 coalesce(col1, crossJoinTBL.Col1Avg) 来完成,但这并不是真正的 UDF。如果您要传入表列并使用动态 SQL 来计算平均值并再次使用可能有效的合并...

标签: java sql scala apache-spark


【解决方案1】:

一般来说这里不需要UDF。你真的只是聚合表:

val df = Seq(
  (Some(2), None, Some(3)), (Some(4), Some(3), Some(3)),
  (Some(6), Some(5), None), (Some(8), None, Some(2)),
  (None, Some(6), Some(4)), (Some(5), Some(2), Some(8))
).toDF("col1", "col2", "col3").alias("df")

val means = df.agg(df.columns.map(c => (c -> "avg")).toMap)

并用coalesce广播笛卡尔:

val exprs = df.columns.map(c => coalesce(col(c), col(s"avg($c)")).alias(c))

df.join(broadcast(means)).select(exprs: _*)

【讨论】:

  • 优秀。这很完美。非常感谢。必须添加以下库。 import sqlctx.implicits._ import org.apache.spark.sql.functions.{coalesce, lit, broadcast}
  • zero323 你的 Scala Spark 技能太疯狂了......不过,如果你能再详细说明一下你的超级漂亮和紧凑的代码是如何工作的,我们将不胜感激。
  • 另外,与我的代码中的所有其他语句相比,df.join(broadcast(means)).select(exprs: _*) 这一行确实需要很长时间。有没有更好的方法来做到这一点?提前致谢。
  • 在 Spark 2.0+ 版本中,将 ``df.join` 替换为 df.crossJoin 以避免出现 org.apache.spark.sql.AnalysisException
猜你喜欢
  • 2021-01-19
  • 1970-01-01
  • 2017-07-10
  • 2021-10-08
  • 2020-03-27
  • 1970-01-01
  • 2018-12-14
  • 2017-02-24
  • 2019-09-17
相关资源
最近更新 更多