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
from pyspark.sql import functions as dbf
@dbf.arrow_udtf(returnType=<returnType>)
class MyUDTF:
def eval(self, ...):
...
parameter
parameter | Type | Description |
|---|---|---|
|
| La classe de gestionnaire de fonction de table définie par l’utilisateur Python. |
|
| 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 :
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 :
@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()
- 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