Cas d'usage PySpark
Cinq cas d'usage PySpark avec le SDK Python de la Data Platform, de la lecture d'un bucket à l'écriture vers un stockage objet
Objectif
Ce guide parcourt cinq cas d'usage PySpark avec le SDK Python de la Data Platform : lire depuis un bucket, interroger des bases de données avec SQL, extraire un dataframe, et écrire vers un bucket ou vers un stockage objet.
Cas d'usage 1 : lire depuis un Bucket et écrire vers un Dataset
Cet exemple montre comment récupérer des données depuis le Bucket de la Data Platform bucket_test et les insérer dans le dataset default_dataset du Lakehouse Manager.
Veuillez consulter le connecteur Buckets de la Data Platform et le connecteur Dataset du Lakehouse Manager pour plus de détails sur les chaînes de connexion.
from logging import getLogger
from forepaas.dwh import connect, update_metas
from pyspark import SparkContext
from pyspark.sql import SQLContext
logger = getLogger(__name__)
cn_source = connect("dwh/bucket_test/chicago_calendar_full.csv")
cn_default = connect("dwh/default_dataset/")
# La fonction extract_dataframe du connecteur compatible Spark retourne un Spark DataFrame
spark_df = cn_source.extract_dataframe()
logger.notice(f"CSV Columns: {list(spark_df.columns)}")
# insert_dataframe utilise également un Spark DataFrame
cn_default.insert_dataframe("chicago_calendar_full_copy", spark_df)
# À la fin, vous pouvez exécuter update_metas() qui mettra à jour les métadonnées de toutes les tables, afin que depuis le lakehouse manager vous voyiez le bon nombre de lignes
update_metas()
Cas d'usage 2 : interroger les bases de données en SQL
Cet exemple met en avant l'interrogation des bases de données via les méthodes get_spark_options() et get_spark_context() de l'objet Connector de la Data Platform.
from logging import getLogger
from forepaas.dwh import connect
from pyspark import SparkContext
from pyspark.sql import SQLContext
logger = getLogger(__name__)
# get_spark_context() et get_spark_options() sont disponibles pour les connecteurs de bases de données (snowflake, postgresql, mysql)
cn_default = connect("dwh/default_dataset/")
sc_default = cn_default.get_spark_context()
so_default = cn_default.get_spark_options()
# Dépend de Snowflake ou PostgreSQL
sql_driver = "net.snowflake.spark.snowflake" # "jdbc" or "net.snowflake.spark.snowflake"
spark_default = SQLContext(sc_default)
sql = "select * from chicago_calendar_full"
spark_df = spark_default.read.format(sql_driver).options(**so_default).option("query", sql).load()
logger.notice(f"SQL Columns: {list(spark_df.columns)}")
Info
Notez que ce cas d'usage fonctionne uniquement avec MySQL, PostgreSQL et Snowflake
Cet exemple montre comment extraire rapidement des Spark DataFrames avec la méthode extract_dataframe() et aussi comment le faire manuellement avec les méthodes get_spark_url(), get_spark_context() et get_spark_session().
from logging import getLogger
from forepaas.dwh import connect
from pyspark import SparkContext
from pyspark.sql import SQLContext
logger = getLogger(__name__)
# Récupère toutes les tables de default_dataset
cn_default = connect("dwh/default_dataset/")
# Si besoin, vous pouvez afficher toutes les tables
# logger.info(cn_default.list())
cn_source = connect("dwh/bucket_test/chicago_calendar_full.csv")
# Redéfinition manuelle des options du fichier extract_dataframe
spark_df = cn_source.extract_dataframe(options)
logger.notice(f"CSV1 Columns: {list(spark_df.columns)}")
# Lecture manuelle depuis le fichier
# get_spark_url(), get_spark_context() et get_spark_session() sont disponibles pour les connecteurs s3 / buckets
spark_session = cn_source.get_spark_session()
# récupération de stations_rides.csv sous le bucket buc_test
url = cn_source.get_spark_url("", "stations_rides.csv", bucket="buc_test")
logger.notice(f"SparkURL: {url}")
# Utilisez format(file_suffix) pour d'autres fichiers, consultez la documentation spark pour plus d'informations
options= {"encoding": "utf-8", "sep": ";", "header": True}
spark_df = spark_session.read.format("csv").options(**options).load(url)
logger.notice(f"CSV2 Columns: {list(spark_df.columns)}")
Cas d'usage 4 : écrire vers un autre bucket
Cet exemple montre comment utiliser la méthode get_spark_url() pour écrire d'un Bucket de la Data Platform vers un autre.
from logging import getLogger
from forepaas.dwh import connect
from pyspark import SparkContext
from pyspark.sql import SQLContext
logger = getLogger(__name__)
cn_source = connect("dwh/bucket_test/stations_rides.csv")
spark_df = cn_source.extract_dataframe()
url_dst = cn_source.get_spark_url("", "stations_rides_copy.csv", bucket="test2")
logger.notice(f"SparkURL Dest: {url_dst}")
spark_df.write.format("csv").options(**options).save(url_dst)
url_dst = cn_source.get_spark_url("", "stations_rides_copy.parquet", bucket="test3")
logger.notice(f"SparkURL Dest: {url_dst}")
spark_df.write.format("parquet").save(url_dst)
Cas d'usage 5 : écrire vers un stockage objet avec insert_dataframe
La fonction insert_dataframe() simplifie les opérations du cas d'usage précédent.
Elle est disponible pour les connecteurs PySpark compatibles de type stockage objet (actuellement Data Platform Buckets, S31 et Azure Blob Storage)
from logging import getLogger
from forepaas.dwh import connect
logger = getLogger(__name__)
cn_source = connect("dwh/bucket_test/chicago_calendar_full.csv")
spark_df = cn_source.extract_dataframe()
# cn_dest: Data Platform Buckets, S3, Azure Blob Storage
cn_dest = connect("dwh/dest_bucket/")
# Insertion avec un type personnalisé
params={"type": "csv"}
# Insertion avec le type par défaut déduit du nom de fichier et des options par défaut de la platform
# Lève une exception si le suffixe de fichier n'est pas dans ["csv", "json", "parquet"]
# Le chemin de destination sera destination, où vous trouverez un fichier .csv sous le chemin configuré de la source
cn_dest.insert_dataframe("destination", spark_df)
# Insertion avec un type personnalisé
params = {"type": "parquet"}
# Lit depuis params.type ; si non fourni et sans suffixe dans le nom de fichier, une exception sera levée
cn_dest.insert_dataframe("destination", spark_df, params)
# Insertion vers un chemin absolu spécifié
params = {"path": "output/destination", "type": "csv"}
# Le fichier de destination sera output/destination.csv
cn_dest.insert_dataframe("", spark_df, params)
# Insertion avec des options personnalisées
params = {"write_options": {"sep": ",", "header": False}, "type": "csv"}
cn_dest.insert_dataframe("destination", spark_df, params)
# Options par défaut actuelles de la platform :
# CSV: {"encoding": "utf-8", "sep": ";", "header": True}
Info
Lors de l'enregistrement d'un PySpark DataFrame vers un système de fichiers comme S3, PySpark crée un dossier au lieu d'un seul fichier car il traite les données de manière distribuée. À l'intérieur de ce dossier, plusieurs fichiers part-*.csv sont générés, chacun représentant une partition du DataFrame, le nombre de fichiers dépendant des partitions du DataFrame. Un fichier _SUCCESS vide est également créé pour indiquer une opération d'écriture réussie.
Aller plus loin
Si vous avez besoin d'une formation ou d'une assistance technique pour la mise en oeuvre de nos solutions, contactez votre commercial ou cliquez sur ce lien pour obtenir un devis et demander une analyse personnalisée de votre projet à nos experts de l’équipe Professional Services.
Posez vos questions, faites-nous part de vos commentaires et interagissez directement avec l’équipe qui développe la Data Platform sur le canal Discord dédié.
Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre centre d'aide.
Rejoignez notre communauté d'utilisateurs.
1 : S3 est une marque déposée appartenant à Amazon Technologies, Inc. Les services de OVHcloud ne sont pas sponsorisés, approuvés, ou affiliés de quelque manière que ce soit.