Aller au contenu principal

arrow_udtf

Crée une fonction de table définie par l’utilisateur native PyArrow (UDTF). Cette fonction fournit une interface native PyArrow pour les UDTF, où la méthode `eval` reçoit des RecordBatches PyArrow ou des tableaux et renvoie un itérateur de Tables PyArrow ou de RecordBatches. Cela permet un véritable calcul vectorisé sans la surcharge de traitement ligne par ligne.

Syntaxe

Python
from pyspark.sql import functions as dbf

@dbf.arrow_udtf(returnType=<returnType>)
class MyUDTF:
def eval(self, ...):
...

parameter

parameter

Type

Description

cls

class, facultatif

La classe de gestionnaire de fonction 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 un objet StructType ou une chaîne de type de structure au format DDL.

parameter

Type

Description

cls

class, facultatif

La classe de gestionnaire de fonction 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 un objet StructType ou une chaîne de type de structure au format DDL.

Exemples

UDTF avec entrée PyArrow RecordBatch :

Python
import pyarrow as pa
from pyspark.sql.functions import arrow_udtf

@arrow_udtf(returnType="x int, y int")
class MyUDTF:
def eval(self, batch: pa.RecordBatch):
# Process the entire batch vectorized
x_array = batch.column('x')
y_array = batch.column('y')
result_table = pa.table({
'x': x_array,
'y': y_array
})
yield result_table

df = spark.range(10).selectExpr("id as x", "id as y")
MyUDTF(df.asTable()).show()

UDTF avec entrées de tableau PyArrow :

Python
@arrow_udtf(returnType="x int, y int")
class MyUDTF2:
def eval(self, x: pa.Array, y: pa.Array):
# Process arrays vectorized
result_table = pa.table({
'x': x,
'y': y
})
yield result_table

MyUDTF2(lit(1), lit(2)).show()
remarque
  • La méthode eval doit accepter les RecordBatches PyArrow ou les tableaux comme entrée.
  • La méthode eval doit renvoyer des tables PyArrow ou des RecordBatches en sortie