Aller au contenu principal

APIs de fonction pandas

Les APIs de fonction pandas vous permettent d'appliquer directement une fonction native Python qui prend et produit des instances pandas à un DataFrame PySpark. De même que les fonctions définies par l'utilisateur pandas, les API de fonction utilisent également Apache Arrow pour transférer les données et pandas pour travailler avec les données ; cependant, les indications de type Python sont facultatives dans les APIs de fonction pandas.

Il existe trois types d'APIs de fonction pandas :

  • Carte groupée
  • Carte
  • Map cogroupée

Les APIs de fonction pandas exploitent la même logique interne que celle utilisée par l’exécution d’UDF pandas. Ils partagent des caractéristiques telles que PyArrow, les types SQL pris en charge et les configurations.

Pour plus d'informations, consultez le billet de blog Nouvelles UDF Pandas et indications de type Python dans la prochaine version d'Apache Spark 3.0.

Carte groupée

Vous transformez vos données groupées à l'aide de groupBy().applyInPandas() pour implémenter le modèle « diviser-appliquer-combiner ». Diviser-appliquer-combiner consiste en trois étapes :

  • Fractionnez les données en groupes en utilisant DataFrame.groupBy.
  • Appliquez une fonction à chaque groupe. L'entrée et la sortie de la fonction sont toutes deux pandas.DataFrame. Les données d'entrée contiennent toutes les lignes et colonnes pour chaque groupe.
  • Combinez les résultats dans une nouvelle DataFrame.

Pour utiliser groupBy().applyInPandas(), vous devez définir ce qui suit :

  • Une fonction Python qui définit le calcul pour chaque groupe
  • Un objet StructType ou une chaîne qui définit le schéma de la sortie DataFrame

Les étiquettes de colonne du pandas.DataFrame retourné doivent soit correspondre aux noms de champ dans le schéma de sortie défini si elles sont spécifiées comme des chaînes de caractères, soit correspondre aux types de données de champ par position si elles ne sont pas des chaînes de caractères, par exemple, des index entiers. Consultez pandas.DataFrame pour savoir comment étiqueter les colonnes lors de la construction d'un pandas.DataFrame.

Toutes les données d'un groupe sont chargées en mémoire avant l'application de la fonction. Cela peut entraîner des exceptions de mémoire insuffisante, en particulier si les tailles de groupe sont asymétriques. La configuration pour maxRecordsPerBatch n'est pas appliquée aux groupes, et il vous incombe de vous assurer que les données groupées tiennent dans la mémoire disponible.

L'exemple suivant montre comment utiliser groupby().apply() pour soustraire la moyenne de chaque valeur du groupe.

Python
df = spark.createDataFrame(
[(1, 1.0), (1, 2.0), (2, 3.0), (2, 5.0), (2, 10.0)],
("id", "v"))

def subtract_mean(pdf):
# pdf is a pandas.DataFrame
v = pdf.v
return pdf.assign(v=v - v.mean())

df.groupby("id").applyInPandas(subtract_mean, schema="id long, v double").show()
# +---+----+
# | id| v|
# +---+----+
# | 1|-0.5|
# | 1| 0.5|
# | 2|-3.0|
# | 2|-1.0|
# | 2| 4.0|
# +---+----+

Pour une utilisation détaillée, consultez pyspark.sql.GroupedData.applyInPandas.

Carte

Vous effectuez des opérations de mappage avec des instances pandas en DataFrame.mapInPandas() afin de transformer un itérateur de pandas.DataFrame en un autre itérateur de pandas.DataFrame qui représente le DataFrame PySpark actuel et renvoie le résultat sous forme de DataFrame PySpark.

La fonction sous-jacente prend en entrée et renvoie un itérateur de pandas.DataFrame. Il peut retourner une sortie de longueur arbitraire, contrairement à certaines UDFs pandas telles que Série à Série.

L'exemple suivant montre comment utiliser mapInPandas():

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, schema=df.schema).show()
# +---+---+
# | id|age|
# +---+---+
# | 1| 21|
# +---+---+

Pour une utilisation détaillée, consultez pyspark.sql.DataFrame.mapInPandas.

Carte cogroupée

Pour les opérations de carte cogroupées avec des instances pandas, utilisez DataFrame.groupby().cogroup().applyInPandas() pour cogrouper deux DataFramePySpark par une clé commune, puis appliquer une fonction Python à chaque cogroupe comme indiqué :

  • Mélangez les données de sorte que les groupes de chaque DataFrame partageant une clé soient cogroupés.
  • Appliquez une fonction à chaque cogroupe. L'entrée de la fonction est de deux pandas.DataFrame (avec un tuple facultatif représentant la clé). La sortie de la fonction est un pandas.DataFrame.
  • Combinez les pandas.DataFramede tous les groupes dans un nouveau DataFrame PySpark.

Pour utiliser groupBy().cogroup().applyInPandas(), vous devez définir ce qui suit :

  • Une fonction Python qui définit le calcul pour chaque cogroupe.
  • Un objet StructType ou une chaîne qui définit le schéma de sortie PySpark DataFrame.

Les étiquettes de colonne du pandas.DataFrame retourné doivent soit correspondre aux noms de champ dans le schéma de sortie défini si elles sont spécifiées comme des chaînes de caractères, soit correspondre aux types de données de champ par position si elles ne sont pas des chaînes de caractères, par exemple, des index entiers. Consultez pandas.DataFrame pour savoir comment étiqueter les colonnes lors de la construction d'un pandas.DataFrame.

Toutes les données d'un co-groupe sont chargées en mémoire avant que la fonction ne soit appliquée. Cela peut entraîner des exceptions de mémoire insuffisante, en particulier si les tailles de groupe sont asymétriques. La configuration pour maxRecordsPerBatch n'est pas appliquée, et il vous incombe de vous assurer que les données co-groupées tiennent dans la mémoire disponible.

L’exemple suivant montre comment utiliser groupby().cogroup().applyInPandas() pour effectuer une asof join entre deux datasets.

Python
import pandas as pd

df1 = spark.createDataFrame(
[(20000101, 1, 1.0), (20000101, 2, 2.0), (20000102, 1, 3.0), (20000102, 2, 4.0)],
("time", "id", "v1"))

df2 = spark.createDataFrame(
[(20000101, 1, "x"), (20000101, 2, "y")],
("time", "id", "v2"))

def asof_join(l, r):
return pd.merge_asof(l, r, on="time", by="id")

df1.groupby("id").cogroup(df2.groupby("id")).applyInPandas(
asof_join, schema="time int, id int, v1 double, v2 string").show()
# +--------+---+---+---+
# | time| id| v1| v2|
# +--------+---+---+---+
# |20000101| 1|1.0| x|
# |20000102| 1|3.0| x|
# |20000101| 2|2.0| y|
# |20000102| 2|4.0| y|
# +--------+---+---+---+

Pour une utilisation détaillée, consultez pyspark.sql.PandasCogroupedOps.applyInPandas.