Digamos que tengo un marco de datos de la siguiente manera:
| id |col | 1 | "A,B,C" | 2 | "D,C" | 3 | "B,C,A" | 4 | Noney el diccionario es:
d = {'A': 1, 'B': 2, 'C': 3, 'D': 4}el marco de datos de salida debe ser:
| id |col | 1 | "A" | 2 | "C" | 3 | "A" | 4 | NoneHigher Order Functions - Transform se puede usar para asociar un rango a los elementos en col según el diccionario y luego ordenarlos para obtener el elemento con el rango más bajo.
from pyspark.sql import functions as F from itertools import chain data = [(1, "A,B,C",), (2, "D,C",), (3, "B,C,A",), (4, None,), ] df = spark.createDataFrame(data, ("id", "col", )) d = {'A': 1, 'B': 2, 'C': 3, 'D': 4} mapper = F.create_map([F.lit(c) for c in chain.from_iterable(d.items())]) """ Mapper has the value Column<'map(A, 1, B, 2, C, 3, D, 4)'> """ (df.withColumn("col", F.split(F.col("col"), ",")) # Split string to create an array .withColumn("mapper", mapper) # Add mapping columing to the dataframe .withColumn("col", F.expr("transform(col, x -> struct(mapper[x] as rank, x as col))")) # Iterate over array and look up rank from mapper .withColumn("col", F.array_min(F.col("col")).col) # array_min find minimum value based on the first struct field ).select("id", "col").show() """ +---+----+ | id| col| +---+----+ | 1| A| | 2| C| | 3| A| | 4|null| +---+----+ """Aquí hay otra solución con orden de estructura como respuesta de @Nithish pero usando arrays_zip y array_min en su lugar:
col ordenados import pyspark.sql.functions as F df = spark.createDataFrame([(1, "A,B,C"), (2, "D,C"), (3, "B,C,A"), (4, None)], ["id", "col"]) d = {'A': 1, 'B': 2, 'C': 3, 'D': 4} result = df.withColumn( "col", F.array_min( F.arrays_zip( F.array(*[F.lit(d[x]) for x in sorted(d)]), F.array_sort(F.split("col", ",")) ) )["1"] ) result.show() #+---+----+ #| id| col| #+---+----+ #| 1| A| #| 2| C| #| 3| A| #| 4|null| #+---+----+Supongo que desea ordenar las letras de acuerdo con los valores dados en el diccionario d .
Entonces, puedes hacer lo siguiente:
from pyspark.sql import Row from pyspark.sql import SparkSession import pyspark.sql.functions as F import pyspark.sql.types as T spark = SparkSession.builder.master("local").appName("sort_column_test").getOrCreate() df = spark.createDataFrame(data=(Row(1, "A,B,C",), Row(2, "D,C",), Row(3, "B,C,A",), Row(4, None)), schema="id:int, col:string") d = {'A': 1, 'B': 2, 'C': 3, 'D': 4} # Define a sort UDF that sorts the array according to the dictionary 'd', also handles None arrays sort_udf = F.udf(lambda array: sorted(array, key=lambda x: d[x]) if array is not None else None, T.ArrayType(T.StringType())) df = df.withColumn("col", sort_udf(F.split(F.col("col"), ",")).getItem(0)) df.show() """ +---+----+ | id| col| +---+----+ | 1| A| | 2| C| | 3| A| | 4|null| +---+----+ """