Spark extrai valores da string e atribui como coluna
Precisa de ajuda na análise de uma string, onde contém valores para cada atributo. abaixo está minha string de amostra ...
Type=<Series VR> Model=<1Ac4> ID=<34> conn seq=<2>
Type=<SeriesX> Model=<12Q3> ID=<231> conn seq=<3423123>
do acima, tenho que gerar as colunas com os valores abaixo.
Type | Model | Id | conn seq
----------------------------
Series VR | 1Ac4 | 34 | 2
SeriesX | 12Q3 | 231 | 3423123
não tenho certeza de como analisá-lo por regex / split e usando withColumn ().
ajuda é apreciada. obrigado
Respostas
Com base em sua amostra, você pode converter a string em mapa usando a função str_to_map do SparkSQL e, em seguida, selecionar valores das chaves de mapa desejadas (o código abaixo assume que o nome da coluna StringType é value):
from pyspark.sql import functions as F
keys = ['Type', 'Model', 'ID', 'conn seq']
df.withColumn("m", F.expr("str_to_map(value, '> *', '=<')")) \
.select("*", *[ F.col('m')[k].alias(k) for k in keys ]) \
.show()
+--------------------+--------------------+---------+-----+---+--------+
| value| m| Type|Model| ID|conn seq|
+--------------------+--------------------+---------+-----+---+--------+
|Type=<Series VR> ...|[Type -> Series V...|Series VR| 1Ac4| 34| 2|
|Type=<SeriesX> Mo...|[Type -> SeriesX,...| SeriesX| 12Q3|231| 3423123|
+--------------------+--------------------+---------+-----+---+--------+
Observações: aqui usamos o padrão regex > *para dividir pares e o padrão =<para dividir chave / valor. Verifique neste link se keysos mapas são dinâmicos e não podem ser predefinidos, apenas certifique-se de filtrar a chave EMPTY.
Editar: Com base em comentários, para fazer pesquisas que não diferenciam maiúsculas de minúsculas nas chaves do mapa. para Spark 2.3 , podemos usar pandas_udf para pré-processar a valuecoluna antes de usar a função str_to_map:
configurar o padrão regex para chaves combinadas (na captura do grupo 1). aqui usamos
(?i)para configurar a correspondência sem distinção entre maiúsculas e minúsculas e adicionar duas âncoras\be(?==), de modo que as subcadeias correspondidas tenham um limite de palavra à esquerda e seguido por uma=marca à direita.ptn = "(?i)\\b({})(?==)".format('|'.join(keys)) print(ptn) #(?i)\b(Type|Model|ID|conn seq)(?==)configure pandas_udf para que possamos usar Series.str.replace () e definir um retorno de chamada ($ 1 minúsculo) como substituição:
lower_keys = F.pandas_udf(lambda s: s.str.replace(ptn, lambda m: m.group(1).lower()), "string")converter todas as chaves correspondentes em minúsculas:
df1 = df.withColumn('value', lower_keys('value')) +-------------------------------------------------------+ |value | +-------------------------------------------------------+ |type=<Series VR> model=<1Ac4> id=<34> conn seq=<2> | |type=<SeriesX> model=<12Q3> id=<231> conn seq=<3423123>| +-------------------------------------------------------+use str_to_map para criar o mapa e, em seguida, use
k.lower()como chaves para encontrar seus valores correspondentes.df1.withColumn("m", F.expr("str_to_map(value, '> *', '=<')")) \ .select("*", *[ F.col('m')[k.lower()].alias(k) for k in keys ]) \ .show()
Observação: caso você possa usar Spark 3.0+ no futuro, pule as etapas acima e use a função transform_keys em seu lugar:
df.withColumn("m", F.expr("str_to_map(value, '> *', '=<')")) \
.withColumn("m", F.expr("transform_keys(m, (k,v) -> lower(k))")) \
.select("*", *[ F.col('m')[k.lower()].alias(k) for k in keys ]) \
.show()
Para Spark 2.4+ , substitua transform_keys(...)pelo seguinte:
map_from_entries(transform(map_keys(m), k -> (lower(k), m[k])))