【问题标题】:Spark udf initializationSpark udf 初始化
【发布时间】:2016-02-13 09:31:02
【问题描述】:

我想在 Spark SQL 中创建一个自定义的基于正则表达式的 UDF。我的偏好是创建一个内存驻留

 Map[String,Pattern]

其中 Pattern 指的是字符串键的编译正则表达式版本。但要做到这一点,我们需要将地图创建放入 UDF 的“初始化”函数中。

那么 Spark udf 是否有任何结构支持跨调用的持久状态(通过 Spark SQL)?

请注意,HIVE 确实支持 UDF 的生命周期。我用它来生成解析树作为初始化的一部分,以便 UDF 的实际调用针对闪电般快速的树,而不涉及解析。

【问题讨论】:

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


    【解决方案1】:

    让我们从导入和一些虚拟数据开始:

    import org.apache.spark.sql.functions.udf
    import scala.util.matching.Regex
    import java.util.regex.Pattern
    
    val df = sc.parallelize(Seq(
      ("foo", "this is bar"), ("foo", "this is foo"),
      ("bar", "foobar"), ("bar", "foo and foo")
    )).toDF("type", "value")
    

    和地图:

    val patterns: Map[String, Pattern] = Seq(("foo", ".*foo.*"), ("bar", ".*bar.*"))
       .map{case (k, v) => (k, new Regex(v).pattern)}
       .toMap
    

    现在我看到了两个不同的选项:

    • 使patterns成为udf内部引用的广播变量

      val patternsBd = sc.broadcast(patterns)
      
      val typeMatchedViaBroadcast = udf((t: String, v: String) =>
        patternsBd.value.get(t).map(m => m.matcher(v).matches))
      
      df.withColumn("match", typeMatchedViaBroadcast($"type", $"value")).show
      
      // +----+-----------+-----+
      // |type|      value|match|
      // +----+-----------+-----+
      // | foo|this is bar|false|
      // | foo|this is foo| true|
      // | bar|     foobar| true|
      // | bar|foo and foo|false|
      // +----+-----------+-----+
      
    • 在闭包内传递地图

      def makeTypeMatchedViaClosure(patterns: Map[String, Pattern]) = udf(
        (t: String, v: String) => patterns.get(t).map(m => m.matcher(v).matches))
      
      val typeMatchedViaClosure = makeTypeMatchedViaClosure(patterns)
      
      df.withColumn("match", typeMatchedViaClosure($"type", $"value")).show
      
      // +----+-----------+-----+
      // |type|      value|match|
      // +----+-----------+-----+
      // | foo|this is bar|false|
      // | foo|this is foo| true|
      // | bar|     foobar| true|
      // | bar|foo and foo|false|
      // +----+-----------+-----+
      

    【讨论】:

    • V 很好的例子。我在考虑注册 Hive UDF。在这种情况下,闭包将不起作用。我忽略了必须在每个 SQLContext 中重新定义 spark udf。
    • @zero323:提供的两种解决方案之间存在相关差异吗? (我主要对表演感兴趣)
    猜你喜欢
    • 1970-01-01
    • 2018-04-14
    • 2016-04-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-08-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多