Aller au contenu principal

udf

Crée une fonction définie par l'utilisateur (UDF).

Syntaxe

Python
import pyspark.sql.functions as sf

# As a decorator
@sf.udf
def function_name(col):
# function body
pass

# As a decorator with return type
@sf.udf(returnType=<returnType>, useArrow=<useArrow>)
def function_name(col):
# function body
pass

# As a function wrapper
sf.udf(f=<function>, returnType=<returnType>, useArrow=<useArrow>)

parameter

parameter

Type

Description

f

function

Facultatif. Fonction Python si utilisée comme fonction autonome.

returnType

pyspark.sql.types.DataType OU str

Facultatif. Type de renvoi de la fonction définie par l'utilisateur. La valeur peut être un objet DataType ou une chaîne de type au format DDL. Default, StringType.

useArrow

bool

Facultatif. S'il faut utiliser Arrow pour optimiser la (dé)sérialisation. Quand elle est définie sur Aucun, la configuration Spark "spark.sql.execution.pythonUDF.arrow.enabled" prend effet.

parameter

Type

Description

f

function

Facultatif. Fonction Python si utilisée comme fonction autonome.

returnType

pyspark.sql.types.DataType OU str

Facultatif. Type de renvoi de la fonction définie par l'utilisateur. La valeur peut être un objet DataType ou une chaîne de type au format DDL. Default, StringType.

useArrow

bool

Facultatif. S'il faut utiliser Arrow pour optimiser la (dé)sérialisation. Quand elle est définie sur Aucun, la configuration Spark "spark.sql.execution.pythonUDF.arrow.enabled" prend effet.

Exemples

Exemple 1 : création d'UDF à l'aide de lambda, de décorateur et de décorateur avec type de retour.

Python
from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf

slen = udf(lambda s: len(s), IntegerType())

@udf
def to_upper(s):
if s is not None:
return s.upper()

@udf(returnType=IntegerType())
def add_one(x):
if x is not None:
return x + 1

df = spark.createDataFrame([(1, "John Doe", 21)], ("id", "name", "age"))
df.select(slen("name").alias("slen(name)"), to_upper("name"), add_one("age")).show()
Output
+----------+--------------+------------+
|slen(name)|to_upper(name)|add_one(age)|
+----------+--------------+------------+
| 8| JOHN DOE| 22|
+----------+--------------+------------+

Exemple 2 : UDF avec arguments nommés.

Python
from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf, col

@udf(returnType=IntegerType())
def calc(a, b):
return a + 10 * b

spark.range(2).select(calc(b=col("id") * 10, a=col("id"))).show()
Output
+-----------------------------+
|calc(b => (id * 10), a => id)|
+-----------------------------+
| 0|
| 101|
+-----------------------------+

**Exemple 3** : UDF vectorisée utilisant les indications de type pandas Series.

Python
from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf, col, PandasUDFType
import pandas as pd

@udf(returnType=IntegerType())
def pd_calc(a: pd.Series, b: pd.Series) -> pd.Series:
return a + 10 * b

pd_calc.evalType == PandasUDFType.SCALAR
spark.range(2).select(pd_calc(b=col("id") * 10, a="id")).show()
Output
+--------------------------------+
|pd_calc(b => (id * 10), a => id)|
+--------------------------------+
| 0|
| 101|
+--------------------------------+

Exemple 4 : UDF vectorisée utilisant des indicateurs de type de tableau PyArrow.

Python
from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf, col, ArrowUDFType
import pyarrow as pa

@udf(returnType=IntegerType())
def pa_calc(a: pa.Array, b: pa.Array) -> pa.Array:
return pa.compute.add(a, pa.compute.multiply(b, 10))

pa_calc.evalType == ArrowUDFType.SCALAR
spark.range(2).select(pa_calc(b=col("id") * 10, a="id")).show()
Output
+--------------------------------+
|pa_calc(b => (id * 10), a => id)|
+--------------------------------+
| 0|
| 101|
+--------------------------------+

Exemple 5 : UDF Python optimisée par Arrow (default depuis Spark 4.2).

Python
from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf

# Arrow optimization is enabled by default since Spark 4.2
@udf(returnType=IntegerType())
def my_udf(x):
return x + 1

# To explicitly disable Arrow optimization and use pickle-based serialization:
@udf(returnType=IntegerType(), useArrow=False)
def legacy_udf(x):
return x + 1

Exemple 6 : Création d'une UDF non déterministe.

Python
from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf
import random

random_udf = udf(lambda: int(random.random() * 100), IntegerType()).asNondeterministic()