【问题标题】:How to refer user defined collection variable in Spark DataFrame SQL如何在 Spark DataFrame SQL 中引用用户定义的集合变量
【发布时间】:2018-11-09 10:27:04
【问题描述】:

我需要允许用户定义不同的命名集合,他们可以在稍后构建 Spark DataFrame SQL 时使用。

我计划为此目的使用 Spark 广播变量,但基于以下 SO 问题 How to refer broadcast variable in Spark DataFrameSQL 看起来这是不可能的

假设作为用户,我通过应用程序 UI 创建了以下集合:

name: countries_dict
values: Seq("Italy", "France", "United States", "Poland", "Spain")

在另一个应用程序 UI(让我们在不同的页面)中,作为用户,我创建了以下 Spark SQL 查询:

SELECT name, phone, country FROM users

我想通过SELECT name, phone, country FROM users WHERE countries in countries_dict过滤记录

因此,例如,现在我可以通过以下方式创建类似的东西:

val countriesDict = Seq("Italy", "France", "United States", "Poland", "Spain")

val inDict = (s: String) => {
  countriesDict.contains(s)
}

spark.udf.register("in_dict", inDict)

然后:

SELECT name, phone, country FROM users WHERE in_dict(country)

但这种方法的最大问题是,countriesDict 是硬编码在代码中的,而不是根据用户在 UI 上的输入动态创建的。

是否可以以某种方式扩展此方法以支持通过应用程序 UI 使用名称和元素动态创建的集合(由用户)?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    在这里使用广播变量并没有什么意义。即使忽略结构问题,调用 udf 的成本也可能会超过广播的好处(尤其是在结构如此小的情况下)。

    内联查询,如果数据很小(使用您喜欢的 SQL 处理库,以避免 SQL 注入的风险):

    SELECT name, phone, country FROM users 
    WHERE country IN ('Italy', 'France', 'United States', 'Poland', 'Spain')
    

    或者只是将输入转换为DataFrame:

    countriesDict.toDF("country").createOrReplaceTempView("countries")
    

    并使用ANTI JOIN,如果数据足够小,则根据广播阈值自动将其提升为广播连接

    SELECT * 
    FROM users LEFT ANTI JOIN countries 
    ON users.country = countries.country
    

    或带有明确的广播提示

    SELECT  /*+ MAPJOIN(countries) */  * 
    FROM users LEFT ANTI JOIN countries 
    ON users.country = countries.country
    

    最后,您可以跳过 SQL 部分并使用 DataFrame API 和 isin

    spark.table("users").where($"country" isin (countriesDict: _*))
    

    或者如果您确实有需要 UDF 的逻辑:

    import org.apache.spark.sql.functions.typedLit
    
    val f = udf((x: String, xs: Seq[String]) => { xs.contains(x) })
    
    spark.table("users").where(f($"country", typedLit(countriesDict)))
    

    【讨论】:

      【解决方案2】:

      当然,我不知道您的应用程序的 UI 等,但有什么反对将集合转换为数据框的吗?当然你不能使用WHERE countries in countries_dict 语法,但你必须使用连接。 但是当连接的数据帧低于某个阈值时,Spark 会自动执行连接作为广播,例如如Mastering Apache Spark中所述

      您只需要一些存储空间,用户可以在其中存储这些小数据帧的内容,例如作为 CSV 文件。

      【讨论】:

        猜你喜欢
        • 2015-11-13
        • 1970-01-01
        • 2014-01-20
        • 1970-01-01
        • 1970-01-01
        • 2012-10-29
        • 1970-01-01
        • 2018-03-19
        相关资源
        最近更新 更多