Scala transforma e divide a coluna de string em coluna de MapType no Dataframe
Eu tenho um DataFrame com uma coluna, contendo urls de solicitação de rastreamento com campos dentro, que se parece com este
df.show(truncate = false)
+--------------------------------
| request_uri
+-----------------------------------
| /i?aid=fptplay&ast=1582163970763&av=4.6.1&did=83295772a8fee349 ...
| /i?p=fplay-ottbox-2019&av=2.0.18&nt=wifi&ov=9&tv=1.0.0&tz=GMT%2B07%3A00 ...
| ...
Eu preciso transformar esta coluna em algo parecido com isto
df.show(truncate = false)
+--------------------------------
| request_uri
+--------------------------------
| (aid -> fptplay, ast -> 1582163970763, tz -> [timezone datatype], nt -> wifi , ...)
| (p -> fplay-ottbox-2019, av -> 2.0.18, ov -> 9, tv -> 1.0.0 , ...)
| ...
Basicamente, preciso dividir os nomes dos campos (delimitador = "&") e seus valores em algum tipo de MapType e adicioná-lo à coluna.
Alguém pode me dar dicas de como escrever uma função personalizada para dividir a coluna da string em uma coluna MapType? Disseram-me para usar withColumn () e mapPartition, mas não sei como implementá-lo de uma maneira que divida as strings e as converta em MapType.
Qualquer ajuda, mesmo que mínima, seria muito apreciada. Sou completamente novo no Scala e estou preso nisso há uma semana.
Respostas
A solução é usar UserDefinedFunctions .
Vamos analisar o problema um passo de cada vez.
// We need a function which converts strings into maps
// based on the format of request uris
def requestUriToMap(s: String): Map[String, String] = {
s.stripPrefix("/i?").split("&").map(elem => {
val pair = elem.split("=")
(pair(0), pair(1)) // evaluate each element to a tuple
}).toMap
}
// Now we convert this function into a UserDefinedFunction.
import org.apache.spark.sql.functions.{col, udf}
// Given to a request uri string, convert it to map, the correct format is assumed.
val requestUriToMapUdf = udf((requestUri: String) => requestUriToMap(requestUri))
Agora, testamos.
// Test data
val df = Seq(
("/i?aid=fptplay&ast=1582163970763&av=4.6.1&did=83295772a8fee349"),
("/i?p=fplay-ottbox-2019&av=2.0.18&nt=wifi&ov=9&tv=1.0.0&tz=GMT%2B07%3A00")
).toDF("request_uri")
df.show(false)
//+-----------------------------------------------------------------------+
//|request_uri |
//+-----------------------------------------------------------------------+
//|/i?aid=fptplay&ast=1582163970763&av=4.6.1&did=83295772a8fee349 |
//|/i?p=fplay-ottbox-2019&av=2.0.18&nt=wifi&ov=9&tv=1.0.0&tz=GMT%2B07%3A00|
//+-----------------------------------------------------------------------+
// Now we execute our UDF to create a column, using the same name replaces that column
val mappedDf = df.withColumn("request_uri", requestUriToMapUdf(col("request_uri")))
mappedDf.show(false)
//+---------------------------------------------------------------------------------------------+
//|request_uri |
//+---------------------------------------------------------------------------------------------+
//|[aid -> fptplay, ast -> 1582163970763, av -> 4.6.1, did -> 83295772a8fee349] |
//|[av -> 2.0.18, ov -> 9, tz -> GMT%2B07%3A00, tv -> 1.0.0, p -> fplay-ottbox-2019, nt -> wifi]|
//+---------------------------------------------------------------------------------------------+
mappedDf.printSchema
//root
// |-- request_uri: map (nullable = true)
// | |-- key: string
// | |-- value: string (valueContainsNull = true)
mappedDf.schema
//org.apache.spark.sql.types.StructType = StructType(StructField(request_uri,MapType(StringType,StringType,true),true))
E era isso que você queria.
Alternativa: Se você não tem certeza se sua string está em conformidade, você pode tentar uma variação diferente da função que tem sucesso mesmo se a string não estiver de acordo com o formato assumido (não contém =ou a entrada é uma String vazia).
def requestUriToMapImproved(s: String): Map[String, String] = {
s.stripPrefix("/i?").split("&").map(elem => {
val pair = elem.split("=")
pair.length match {
case 0 => ("", "") // in case the given string produces an array with no elements e.g. "=".split("=") == Array()
case 1 => (pair(0), "") // in case the given string contains no = and produces a single element e.g. "potato".split("=") == Array("potato")
case _ => (pair(0), pair(1)) // normal case e.g. "potato=masher".split("=") == Array("potato", "masher")
}
}).toMap
}
O código a seguir executa um processo de divisão em duas fases. Porque uris não tem uma estrutura específica, você pode fazer isso com uma UDF como:
val keys = List("aid", "p", "ast", "av", "did", "nt", "ov", "tv", "tz")
def convertToMap(keys: List[String]) = udf {
(in : mutable.WrappedArray[String]) =>
in.foldLeft[Map[String, String]](Map()){ (a, str) =>
keys.flatMap { key =>
val regex = s"""${key}=""" val arr = str.split(regex) val value = { if(arr.length == 2) arr(1) else "" } if(!value.isEmpty) a + (key -> value) else a }.toMap } } df.withColumn("_tmp", split($"request_uri","""((&)|(\?))"""))
.withColumn("map_result", convertToMap(keys)($"_tmp")) .select($"map_result")
.show(false)
fornece uma coluna MapType:
+------------------------------------------------------------------------------------------------+
|map_result |
+------------------------------------------------------------------------------------------------+
|Map(aid -> fptplay, ast -> 1582163970763, av -> 4.6.1, did -> 83295772a8fee349) |
|Map(av -> 2.0.18, ov -> 9, tz -> GMT%2B07%3A00, tv -> 1.0.0, p -> fplay-ottbox-2019, nt -> wifi)|
+------------------------------------------------------------------------------------------------+