【问题标题】:Spark scala data frame udf returning rowsSpark scala数据帧udf返回行
【发布时间】:2018-09-17 17:42:30
【问题描述】:

假设我有一个数据框,其中包含一个列(称为 colA),它是行的序列。我想为 colA 的每条记录附加一个新字段。 (而且新的归档是和之前的记录相关联的,所以我要写一个udf。) 这个udf应该怎么写?

我尝试编写一个 udf,它将 colA 作为输入,并输出 Seq[Row],其中每条记录都包含新文件。但问题是 udf 无法返回 Seq[Row]/ 例外是“不支持 org.apache.spark.sql.Row 类型的架构”。 我该怎么办?

我写的udf: val convert = udf[Seq[Row], Seq[Row]](blablabla...) 例外是 java.lang.UnsupportedOperationException:不支持 org.apache.spark.sql.Row 类型的架构

【问题讨论】:

  • 如果您的最终组合行是固定的,则创建一个案例类并使用它。您必须向我们提供更多信息以获得详细答案
  • 输入列的类型是什么?即你的“行”中有什么?

标签: scala apache-spark user-defined-functions


【解决方案1】:

从 spark 2.0 开始,您可以创建返回 Row / Seq[Row] 的 UDF,但您必须提供返回类型的架构,例如如果您使用双精度数组:

val schema = ArrayType(DoubleType)

val myUDF = udf((s: Seq[Row]) => {
  s // just pass data without modification
}, schema)

但我真的无法想象这在哪里有用,我宁愿从 UDF 中返回元组或案例类(或其 Seq)。

编辑:如果您的行包含超过 22 个字段(元组/案例类的字段限制),这可能会很有用

【讨论】:

  • 太棒了!这就是我需要的。情况是我希望新列使用 df.withColumn() 覆盖原始列。你能告诉我spark中case类的基本用法吗?
  • 这是一个 java udf。 scala 的正确格式是什么?
  • @kevinno 这是 Scala
猜你喜欢
  • 1970-01-01
  • 2017-05-21
  • 1970-01-01
  • 2017-09-16
  • 2017-11-01
  • 2020-08-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多