Aller au contenu principal

mapInArrow

Mappe un itérateur de batchs dans le DataFrame actuel à l'aide d'une fonction native Python qui est exécutée sur des pyarrow.RecordBatchs à la fois en entrée et en sortie, et renvoie le résultat sous forme de DataFrame.

Syntaxe

mapInArrow(func: "ArrowMapIterFunction", schema: Union[StructType, str], barrier: bool = False, profile: Optional[ResourceProfile] = None)

parameter

parameter

Type

Description

func

fonction

une fonction native Python qui prend un itérateur de pyarrow.RecordBatch, et produit un itérateur de pyarrow.RecordBatch.

schema

DataType ou str

le type de retour du func dans PySpark. La valeur peut être un objet pyspark.sql.types.DataType ou une chaîne de type au format DDL.

barrier

bool, facultatif, default : false

Utilisez l'exécution en mode barrière, en vous assurant que tous les Worker Python de l'étape seront lancés simultanément.

profile

ResourceProfile, facultatif

Le ResourceProfile facultatif à utiliser pour mapInArrow.

parameter

Type

Description

func

fonction

une fonction native Python qui prend un itérateur de pyarrow.RecordBatch, et produit un itérateur de pyarrow.RecordBatch.

schema

DataType ou str

le type de retour du func dans PySpark. La valeur peut être un objet pyspark.sql.types.DataType ou une chaîne de type au format DDL.

barrier

bool, facultatif, default : false

Utilisez l'exécution en mode barrière, en vous assurant que tous les Worker Python de l'étape seront lancés simultanément.

profile

ResourceProfile, facultatif

Le ResourceProfile facultatif à utiliser pour mapInArrow.

Renvoie

DataFrame

Exemples

Python
import pyarrow as pa
df = spark.createDataFrame([(1, 21), (2, 30)], ("id", "age"))
def filter_func(iterator):
for batch in iterator:
pdf = batch.to_pandas()
yield pa.RecordBatch.from_pandas(pdf[pdf.id == 1])
df.mapInArrow(filter_func, df.schema).show()
# +---+---+
# | id|age|
# +---+---+
# | 1| 21|
# +---+---+

df.mapInArrow(filter_func, df.schema, barrier=True).collect()
# [Row(id=1, age=21)]