Aller au contenu principal

Conversion entre les DataFrames PySpark et pandas

Découvrez comment convertir les DataFrames Apache Spark vers et depuis les DataFrames Pandas à l’aide d’Apache Arrow dans Databricks.

Apache Arrow et PyArrow

Apache Arrow est un format de données colonne en mémoire utilisé dans Apache Spark pour transférer efficacement les données entre les processus JVM et Python. Ceci est bénéfique pour les développeurs Python qui travaillent avec des données pandas et NumPy. Cependant, son utilisation nécessite des modifications mineures de configuration ou de code pour assurer la compatibilité et en tirer le meilleur parti.

PyArrow est une liaison Python pour Apache Arrow et est installé dans Databricks Runtime. Pour plus d'informations sur la version de PyArrow disponible dans chaque version de Databricks Runtime, consultez les notes de publication de Databricks Runtime versions et compatibilité.

Types SQL pris en charge

Tous les types de données Spark SQL sont pris en charge par la conversion basée sur Arrow, à l'exception de ArrayType de TimestampType. MapType et ArrayType de StructType imbriqués sont pris en charge uniquement lors de l'utilisation de PyArrow 2.0.0 et des versions ultérieures. StructType est représenté comme un pandas.DataFrame au lieu de pandas.Series.

Convertir les DataFrames PySpark vers et depuis les DataFrames pandas

Arrow est disponible comme optimisation lors de la conversion d'un PySpark DataFrame en un DataFrame pandas avec toPandas() et lors de la création d'un PySpark DataFrame à partir d'un DataFrame pandas avec createDataFrame(pandas_df).

Pour utiliser Arrow pour ces méthodes, définissez la configuration Spark spark.sql.execution.arrow.pyspark.enabled sur true. Cette configuration est activée par default, à l'exception des clusters high concurrency ainsi que des clusters d'isolation d'utilisateur dans les workspaces sur lesquels Unity Catalog est activé.

En outre, les optimisations activées par spark.sql.execution.arrow.pyspark.enabled pourraient revenir à une implémentation non-Arrow si une erreur se produit avant le calcul dans Spark. Vous pouvez contrôler ce comportement à l'aide de la configuration Spark spark.sql.execution.arrow.pyspark.fallback.enabled.

Exemple

Python
import numpy as np
import pandas as pd

# Enable Arrow-based columnar data transfers
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")

# Generate a pandas DataFrame
pdf = pd.DataFrame(np.random.rand(100, 3))

# Create a Spark DataFrame from a pandas DataFrame using Arrow
df = spark.createDataFrame(pdf)

# Convert the Spark DataFrame back to a pandas DataFrame using Arrow
result_pdf = df.select("*").toPandas()

L'utilisation des optimisations Arrow produit les mêmes résultats que lorsque Arrow n'est pas activé. Même avec Arrow, toPandas() entraîne la collecte de tous les enregistrements du DataFrame dans le programme Driver et doit être effectuée sur un petit sous-ensemble des données.

De plus, tous les types de données Spark ne sont pas pris en charge et une erreur peut survenir si une colonne a un type non pris en charge. Si une erreur survient pendant createDataFrame(), Spark crée le DataFrame sans Arrow.