【问题标题】:How to pass whole Row to UDF - Spark DataFrame filter如何将整行传递给 UDF - Spark DataFrame 过滤器
【发布时间】:2015-10-27 07:29:40
【问题描述】:

我正在为具有大量内部结构的复杂 JSON 数据集编写过滤器函数。传递单个列太麻烦了。

所以我声明了以下 UDF:

val records:DataFrame = = sqlContext.jsonFile("...")
def myFilterFunction(r:Row):Boolean=???
sqlc.udf.register("myFilter", (r:Row)=>myFilterFunction(r))

直觉上我认为它会像这样工作:

records.filter("myFilter(*)=true")

实际的语法是什么?

【问题讨论】:

  • 您能再详细说明一下您的过滤器功能吗?使用Row 会抛弃DataFrame 为您所做的许多优化。
  • 过滤器相当复杂。记录的结构是几个 Map 字段,里面有一堆键值对

标签: apache-spark


【解决方案1】:

如果你想对整行采取一个动作并以分布式的方式处理它,将DataFrame中的行作为结构发送到一个函数,然后转换为字典来执行特定的动作,非常在最终的 DataFrame 上执行 collect 方法很重要,因为 Spark 已激活 LazyLoad 并且不使用完整数据,除非您明确告诉它。

在我的情况下,我应该将 DataFrame 的行作为 Dictionary 对象发送到索引:

  1. 导入库。
  2. 声明 udf 并且 lambda 必须接收行结构。
  3. 执行特定函数,在这种情况下发送到索引字典(行结构转换为字典)。
  4. DataFrame 源执行一个withColum 方法,指示Spark 在每一行中执行此操作,在调用collect 之前,这允许以可分发的方式执行该函数。不要忘记发送到其他 DataFrame 变量。
  5. 执行collect方法,执行进程,分发函数。
from pyspark.sql.functions import udf, struct
from pyspark.sql.types import IntegerType

myUdf = udf(lambda row: sendToES(row.asDict()), IntegerType())
dfWithControlCol = df.withColumn("control_col", myUdf(struct([df[x] for x in df.columns])))
dfWithControlCol.collect()

【讨论】:

    【解决方案2】:

    除了第一个答案。当我们希望所有列都传递给 UDF 时,我们可以使用

     struct("*")
    

    【讨论】:

    • 请添加更多详细信息以扩展您的答案,例如工作代码或文档引用。
    【解决方案3】:
    scala> inputDF
    res40: org.apache.spark.sql.DataFrame = [email: string, first_name: string ... 3 more fields]
    
    scala> inputDF.printSchema
    root
     |-- email: string (nullable = true)
     |-- first_name: string (nullable = true)
     |-- gender: string (nullable = true)
     |-- id: long (nullable = true)
     |-- last_name: string (nullable = true)
    

    现在,我想根据性别字段过滤行。我可以通过使用.filter($"gender" === "Male") 来实现这一点,但我想使用.filter(function)

    所以,定义了我的匿名函数

    val isMaleRow = (r:Row) => {r.getAs("gender") == "Male"}
    
    val isFemaleRow = (r:Row) => { r.getAs("gender") == "Female" }
    
    inputDF.filter(isMaleRow).show()
    
    inputDF.filter(isFemaleRow).show()
    

    我觉得可以以更好的方式完成要求,即无需声明为 UDF 并调用它。

    【讨论】:

    • 能把格式中的非代码解释去掉吗?
    【解决方案4】:

    在调用函数时,您必须使用struct() 函数来构造行,请按照以下步骤操作。

    导入行,

    import org.apache.spark.sql._
    

    定义 UDF

    def myFilterFunction(r:Row) = {r.get(0)==r.get(1)} 
    

    注册 UDF

    sqlContext.udf.register("myFilterFunction", myFilterFunction _)
    

    创建数据帧

    val records = sqlContext.createDataFrame(Seq(("sachin", "sachin"), ("aggarwal", "aggarwal1"))).toDF("text", "text2")
    

    使用 UDF

    records.filter(callUdf("myFilterFunction",struct($"text",$"text2"))).show
    

    当您希望将所有列传递给 UDF 时。

    records.filter(callUdf("myFilterFunction",struct(records.columns.map(records(_)) : _*))).show 
    

    结果:

    +------+------+
    |  text| text2|
    +------+------+
    |sachin|sachin|
    +------+------+
    

    【讨论】:

    • 这有帮助。不是 100% 我需要的,但比替代品更好。
    • 请详细说明我错过了什么以使其 100%,我会尝试根据您的要求更新我的建议。
    • 理想的情况是不必列出列,而是以某种方式引用整行。
    • (很晚了,我相信你现在已经有了答案,但是:struct(df.columns.map(df(_)) : _*)
    • 我可以编辑它,但我只能使用 import org.apache.spark.sql.functions.callUDF。我注意到它是 callUDF。多年来,情况可能发生了变化。
    猜你喜欢
    • 2018-05-10
    • 2023-03-06
    • 2018-02-02
    • 1970-01-01
    • 2019-02-02
    • 1970-01-01
    • 1970-01-01
    • 2021-12-06
    • 2016-06-03
    相关资源
    最近更新 更多