【问题标题】:Scala UDF function to operate on Array Column and return custom valueScala UDF 函数对数组列进行操作并返回自定义值
【发布时间】:2020-06-27 18:17:32
【问题描述】:

我有 2 个这样的数据集

val jsonStr ="""{
    
        "TransactionId": 1,
        "TransactionName": "Name",
        "Order": 12, 
        "ReplaceStrings": [
            "UNDEFINED","INVALID"
        ],
        "Country" : "China"           
     
}"""

val configurations = spark.read.json(Seq(jsonStr).toDS) 

这有我所有的配置和过滤器

My Data
val data =  Seq((1,"Mindy","Devaney","mdevaney0@cnbc.com","Female","United States","UTF-8"),(2,"Charmain","Clear","candriolli1@miitbeian.gov.cn","Female","**China**","UTF-8"),(3,"Dilan","**UNDEFINED**","dphilipeaux2@jalbum.net","Male","**China**","Windows-1252")).toDF("id","Fname","LName","mailid","Gender","Country","Codepage" )

现在,我的要求是加入具有过滤器的配置数据,并使用上述数据检索相应的结果,因为过滤器应用于中国,所有具有 UNDEFINED 作为值的 LName 都将替换为空字符串。

我尝试放置一些 UDF 将其定义为函数,但坚持如何发送 json 值,它是一个包装数组或尝试使用 Seq 数据类型

如果有人看过类似的案例或想法,请与我分享。

【问题讨论】:

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


    【解决方案1】:

    检查下面的代码。

    scala> data.show(false)
    +---+--------+-------------+----------------------------+------+-------------+------------+
    |id |Fname   |LName        |mailid                      |Gender|Country      |Codepage    |
    +---+--------+-------------+----------------------------+------+-------------+------------+
    |1  |Mindy   |Devaney      |mdevaney0@cnbc.com          |Female|United States|UTF-8       |
    |2  |Charmain|Clear        |candriolli1@miitbeian.gov.cn|Female|**China**    |UTF-8       |
    |3  |Dilan   |**UNDEFINED**|dphilipeaux2@jalbum.net     |Male  |**China**    |Windows-1252|
    +---+--------+-------------+----------------------------+------+-------------+------------+
    
    scala> configurations.show(false)
    +-------+-----+--------------------+-------------+---------------+
    |Country|Order|ReplaceStrings      |TransactionId|TransactionName|
    +-------+-----+--------------------+-------------+---------------+
    |China  |12   |[UNDEFINED, INVALID]|1            |Name           |
    +-------+-----+--------------------+-------------+---------------+
    
    
    scala> val check = udf((lname:String,replaceStrings:Seq[String]) => if(replaceStrings.map(d => s"**${d}**").contains(lname)) "" else lname )
    
    scala> data.join(configurations,data("Country").contains(configurations("Country")),"inner").withColumn("LName",check($"LName",$"ReplaceStrings")).drop(configurations("Country")).show(false)
    +---+--------+-----+----------------------------+------+---------+------------+-----+--------------------+-------------+---------------+
    |id |Fname   |LName|mailid                      |Gender|Country  |Codepage    |Order|ReplaceStrings      |TransactionId|TransactionName|
    +---+--------+-----+----------------------------+------+---------+------------+-----+--------------------+-------------+---------------+
    |2  |Charmain|Clear|candriolli1@miitbeian.gov.cn|Female|**China**|UTF-8       |12   |[UNDEFINED, INVALID]|1            |Name           |
    |3  |Dilan   |     |dphilipeaux2@jalbum.net     |Male  |**China**|Windows-1252|12   |[UNDEFINED, INVALID]|1            |Name           |
    +---+--------+-----+----------------------------+------+---------+------------+-----+--------------------+-------------+---------------+
    
    

    【讨论】:

      猜你喜欢
      • 2018-05-13
      • 2019-06-23
      • 2017-11-01
      • 2018-02-17
      • 1970-01-01
      • 1970-01-01
      • 2018-09-17
      • 2023-01-25
      • 1970-01-01
      相关资源
      最近更新 更多