【问题标题】:Spark function aliases - performant udfsSpark 函数别名 - 高性能 udfs
【发布时间】:2019-12-05 14:13:47
【问题描述】:

上下文

在我编写的许多 sql 查询中,我发现自己以完全相同的方式组合 spark 预定义函数,这经常导致冗长和重复代码,而我的开发人员本能就是想要重构它。

所以,我的问题是:是否有某种方法可以为函数组合定义某种 别名 而无需诉诸 udfs(出于性能原因,这是为了避免) - 目标是使代码更清晰,更干净。本质上,我想要的是 udfs 之类的东西,但没有性能损失。此外,这些函数必须可以在可用于spark.sql 调用的 spark-sql 查询中调用。

示例

例如,假设我的业务逻辑是反转一些字符串并像这样对其进行哈希处理:(请注意,这里的函数组合是无关紧要的,重要的是它是现有预定义 spark 函数的某种组合 -可能很多)

SELECT 
    sha1(reverse(person.name)),
    sha1(reverse(person.some_information)),
    sha1(reverse(person.some_other_information))
    ...
FROM person

有没有一种方法可以声明business 函数而无需支付使用udf 的性能代价,允许将上面的代码重写为:

SELECT 
    business(person.name),
    business(person.some_information),
    business(person.some_other_information)
    ...
FROM person

我在 spark 文档和这个网站上搜索了很多,但没有找到实现这一点的方法,这对我来说很奇怪,因为它看起来很自然的需要,我不明白为什么您必须为定义和调用 udf 付出黑盒的代价。

【问题讨论】:

  • 我想你可能已经回答了你自己的问题。

标签: apache-spark apache-spark-sql apache-spark-2.0


【解决方案1】:

有没有一种方法可以在不支付使用 udf 的性能代价的情况下声明业务功能

您不必使用udf,您可以扩展Expression 类,或者对于最简单的操作-UnaryExpression。然后你将不得不实现几个方法,我们开始吧。它原生集成到 Spark 中,此外还可以使用代码生成等一些优势功能。

在您的情况下,添加 business 函数非常简单:

def business(column: Column): Column = {
  sha1(reverse(column))
}

必须可以从可用于 spark.sql 调用的 spark-sql 查询中调用

这比较棘手,但可以实现。
您需要创建自定义函数注册器:

import org.apache.spark.sql.catalyst.FunctionIdentifier
import org.apache.spark.sql.catalyst.expressions.Expression 

object FunctionAliasRegistrar {

val funcs: mutable.Map[String, Seq[Column] => Column] = mutable.Map.empty

  def add(name: String, builder: Seq[Column] => Column): this.type = {
    funcs += name -> builder
    this
  }

  def registerAll(spark: SparkSession) = {
    funcs.foreach { case (alias, builder) => {
      def b(children: Seq[Expression]) = builder.apply(children.map(expr => new Column(expr))).expr
      spark.sessionState.functionRegistry.registerFunction(FunctionIdentifier(alias), b)
    }}
  }
}

那么你可以如下使用它:

FunctionAliasRegistrar
  .add("business1", child => lower(reverse(child.head)))
  .add("business2", child => upper(reverse(child.head)))
  .registerAll(spark) 

dataset.createTempView("data")

spark.sql(
  """
    | SELECT business1(name), business2(name) FROM data
    |""".stripMargin)
.show(false)

输出:

+--------------------+--------------------+
|lower(reverse(name))|upper(reverse(name))|
+--------------------+--------------------+
|sined               |SINED               |
|taram               |TARAM               |
|1taram              |1TARAM              |
|2taram              |2TARAM              |
+--------------------+--------------------+

希望这会有所帮助。

【讨论】:

  • 这确实是我所需要的。我刚刚测试了它(spark 2.4.2),看起来如果你愿意使用 Column 类的构造函数而不是在registerAll 方法中使用Column.apply,那么你不需要把@ 987654331@对象在org.apache.spark.sql! (也许用这个改变来更新答案会很好,因为它让整个事情看起来不像是一个黑客,更像是一个干净的扩展!非常感谢:)
  • 没注意到,更新答案,谢谢。
  • 第二部分不错。
  • 你知道是否可以使用 pyspark 执行此操作,我有一个返回列表达式的 python 函数,我希望能够在 sql 语法和 pyspark 中使用必须创建两个函数。谢谢
猜你喜欢
  • 1970-01-01
  • 2018-09-08
  • 2016-11-12
  • 2020-09-09
  • 2013-06-24
  • 2019-11-23
  • 1970-01-01
  • 2016-06-21
  • 1970-01-01
相关资源
最近更新 更多