Aller au contenu principal

udtf

Crée une fonction de table définie par l'utilisateur (UDTF).

Syntaxe

Python
import pyspark.sql.functions as sf

# As a decorator
@sf.udtf(returnType=<returnType>, useArrow=<useArrow>)
class FunctionClass:
def eval(self, *args):
# function body
yield row_data

# As a function wrapper
sf.udtf(cls=<class>, returnType=<returnType>, useArrow=<useArrow>)

parameter

parameter

Type

Description

cls

class

Facultatif. La classe de gestionnaire de fonctions de table définie par l'utilisateur Python.

returnType

pyspark.sql.types.StructType OU str

Facultatif. Le type de retour de la fonction de table définie par l'utilisateur. La valeur peut être soit un objet StructType, soit une chaîne de caractères de type de structure formatée DDL. Si aucun, la classe du gestionnaire doit fournir la méthode statique analyze.

useArrow

bool

Facultatif. Indique s’il faut utiliser Arrow pour optimiser les (dé)sérialisations. Lorsqu'il est défini sur Aucun, la configuration Spark "spark.sql.execution.pythonUDTF.arrow.enabled" est utilisée.

parameter

Type

Description

cls

class

Facultatif. La classe de gestionnaire de fonctions de table définie par l'utilisateur Python.

returnType

pyspark.sql.types.StructType OU str

Facultatif. Le type de retour de la fonction de table définie par l'utilisateur. La valeur peut être soit un objet StructType, soit une chaîne de caractères de type de structure formatée DDL. Si aucun, la classe du gestionnaire doit fournir la méthode statique analyze.

useArrow

bool

Facultatif. Indique s’il faut utiliser Arrow pour optimiser les (dé)sérialisations. Lorsqu'il est défini sur Aucun, la configuration Spark "spark.sql.execution.pythonUDTF.arrow.enabled" est utilisée.

Exemples

Exemple 1 : Implémentation UDTF de base.

Python
from pyspark.sql.functions import udtf

class TestUDTF:
def eval(self, *args):
yield "hello", "world"

test_udtf = udtf(TestUDTF, returnType="c1: string, c2: string")
test_udtf().show()
Output
+-----+-----+
| c1| c2|
+-----+-----+
|hello|world|
+-----+-----+

**Exemple 2** : UDTF utilisant la syntaxe du décorateur.

Python
from pyspark.sql.functions import udtf, lit

@udtf(returnType="c1: int, c2: int")
class PlusOne:
def eval(self, x: int):
yield x, x + 1

PlusOne(lit(1)).show()
Output
+---+---+
| c1| c2|
+---+---+
| 1| 2|
+---+---+

Exemple 3 : UDTF avec la méthode d'analyse statique.

Python
from pyspark.sql.functions import udtf, lit
from pyspark.sql.types import StructType
from pyspark.sql.udtf import AnalyzeArgument, AnalyzeResult

@udtf
class TestUDTFWithAnalyze:
@staticmethod
def analyze(a: AnalyzeArgument, b: AnalyzeArgument) -> AnalyzeResult:
return AnalyzeResult(StructType().add("a", a.dataType).add("b", b.dataType))

def eval(self, a, b):
yield a, b

TestUDTFWithAnalyze(lit(1), lit("x")).show()
Output
+---+---+
| a| b|
+---+---+
| 1| x|
+---+---+

Exemple 4 : UDTF avec des arguments de mot clé.

Python
from pyspark.sql.functions import udtf, lit
from pyspark.sql.types import StructType
from pyspark.sql.udtf import AnalyzeArgument, AnalyzeResult

@udtf
class TestUDTFWithKwargs:
@staticmethod
def analyze(
a: AnalyzeArgument, b: AnalyzeArgument, **kwargs: AnalyzeArgument
) -> AnalyzeResult:
return AnalyzeResult(
StructType().add("a", a.dataType)
.add("b", b.dataType)
.add("x", kwargs["x"].dataType)
)

def eval(self, a, b, **kwargs):
yield a, b, kwargs["x"]

TestUDTFWithKwargs(lit(1), x=lit("x"), b=lit("b")).show()
Output
+---+---+---+
| a| b| x|
+---+---+---+
| 1| b| x|
+---+---+---+

Exemple 5 : UDTF enregistrée et appelée via SQL.

Python
from pyspark.sql.functions import udtf, lit
from pyspark.sql.types import StructType
from pyspark.sql.udtf import AnalyzeArgument, AnalyzeResult

@udtf
class TestUDTFWithKwargs:
@staticmethod
def analyze(
a: AnalyzeArgument, b: AnalyzeArgument, **kwargs: AnalyzeArgument
) -> AnalyzeResult:
return AnalyzeResult(
StructType().add("a", a.dataType)
.add("b", b.dataType)
.add("x", kwargs["x"].dataType)
)

def eval(self, a, b, **kwargs):
yield a, b, kwargs["x"]

_ = spark.udtf.register("test_udtf", TestUDTFWithKwargs)
spark.sql("SELECT * FROM test_udtf(1, x => 'x', b => 'b')").show()
Output
+---+---+---+
| a| b| x|
+---+---+---+
| 1| b| x|
+---+---+---+

Exemple 6 : UDTF avec optimisation Arrow activée.

Python
from pyspark.sql.functions import udtf, lit

@udtf(returnType="c1: int, c2: int", useArrow=True)
class ArrowPlusOne:
def eval(self, x: int):
yield x, x + 1

ArrowPlusOne(lit(1)).show()
Output
+---+---+
| c1| c2|
+---+---+
| 1| 2|
+---+---+

Exemple 7 : Création d’une UDTF déterministe.

Python
from pyspark.sql.functions import udtf

class PlusOne:
def eval(self, a: int):
yield a + 1,

plus_one = udtf(PlusOne, returnType="r: int").asDeterministic()