Aller au contenu principal

mapInPandas

Mappe un itérateur de lots dans le DataFrame actuel à l'aide d'une fonction native Python qui s'exécute sur des DataFrames pandas à la fois en entrée et en sortie, et renvoie le résultat sous forme de DataFrame.

Syntaxe

mapInPandas(func: "PandasMapIterFunction", 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 pandas.DataFrame, et produit un itérateur de pandas.DataFrame.

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 mapInPandas.

parameter

Type

Description

func

fonction

une fonction native Python qui prend un itérateur de pandas.DataFrame, et produit un itérateur de pandas.DataFrame.

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 mapInPandas.

Renvoie

DataFrame

Exemples

Python
df = spark.createDataFrame([(1, 21), (2, 30)], ("id", "age"))

def filter_func(iterator):
for pdf in iterator:
yield pdf[pdf.id == 1]

df.mapInPandas(filter_func, df.schema).show()
# +---+---+
# | id|age|
# +---+---+
# | 1| 21|
# +---+---+

def mean_age(iterator):
for pdf in iterator:
yield pdf.groupby("id").mean().reset_index()

df.mapInPandas(mean_age, "id: bigint, age: double").show()
# +---+----+
# | id| age|
# +---+----+
# | 1|21.0|
# | 2|30.0|
# +---+----+

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