1. Présentation de Serverless pour Apache Spark
Managed Service pour Apache Spark est un service entièrement géré et hautement évolutif qui permet d'exécuter Apache Spark, Apache Flink, Presto et de nombreux autres outils et frameworks Open Source. Utilisez Managed Apache Spark pour moderniser vos lacs de données, effectuer des tâches d'ETL / ELT et sécuriser la data science à l'échelle mondiale. Managed Apache Spark est également entièrement intégré à plusieurs services Google Cloud, y compris BigQuery, Cloud Storage, Gemini Enterprise Agent Engine et Knowledge Catalog.
Managed Apache Spark est disponible en deux versions :
- Managed Apache Spark sans serveur vous permet d'exécuter des jobs PySpark sans avoir à configurer l'infrastructure ni l'autoscaling. Managed Apache Spark est compatible avec les charges de travail par lot et les sessions / notebooks PySpark.
- Les clusters Apache Spark gérés vous permettent de gérer un cluster Hadoop YARN pour les charges de travail Spark basées sur YARN, en plus des outils Open Source tels que Flink et Presto. Vous pouvez personnaliser vos clusters cloud avec autant de scaling vertical ou horizontal que vous le souhaitez, y compris l'autoscaling.
Dans cet atelier de programmation, vous allez découvrir plusieurs façons d'utiliser Dataproc sans serveur.
À l'origine, Apache Spark a été conçu pour s'exécuter sur des clusters Hadoop, et utilisait YARN comme gestionnaire de ressources. La maintenance de clusters Hadoop requiert un ensemble spécifique d'expertise et la garantie que de nombreux paramètres des clusters sont correctement configurés. En plus de cet ensemble de paramètres, Spark exige également que l'utilisateur en définisse un autre. Les développeurs passent donc plus de temps à configurer leur infrastructure qu'à travailler sur le code Spark lui-même.
Dataproc Serverless élimine la nécessité de configurer manuellement les clusters Hadoop ou Spark. Dataproc Serverless ne s'exécute pas sur Hadoop et utilise sa propre allocation dynamique des ressources pour déterminer ses besoins en ressources, y compris l'autoscaling. Un petit sous-ensemble des propriétés Spark est toujours personnalisable avec Dataproc Serverless, mais dans la plupart des cas, vous n'aurez pas besoin de les modifier.
2. Configurer
Vous allez commencer par configurer votre environnement et les ressources utilisées dans cet atelier de programmation.
Créez un projet Google Cloud. Vous pouvez en utiliser un existant.
Ouvrez Cloud Shell en cliquant dessus dans la barre d'outils de la console Cloud.

Cloud Shell fournit un environnement de shell prêt à l'emploi que vous pouvez utiliser pour cet atelier de programmation.

Cloud Shell définit le nom de votre projet par défaut. Vérifiez-le en exécutant echo $GOOGLE_CLOUD_PROJECT. Si l'ID de votre projet ne s'affiche pas dans le résultat, définissez-le.
export GOOGLE_CLOUD_PROJECT=<your-project-id>
Définissez une région Compute Engine pour vos ressources, par exemple us-central1 ou europe-west2.
export REGION=<your-region>
Activer les API
Cet atelier de programmation utilise les API suivantes :
- BigQuery
- Dataproc
Activez les API nécessaires. Cette opération prend environ une minute. Un message de réussite s'affiche une fois l'opération terminée.
gcloud services enable bigquery.googleapis.com gcloud services enable dataproc.googleapis.com
Configurer l'accès au réseau
Dataproc sans serveur nécessite l'activation de l'accès privé à Google dans la région où vous exécuterez vos jobs Spark, car les pilotes et exécuteurs Spark ne disposent que d'adresses IP privées. Exécutez la commande suivante pour l'activer dans le sous-réseau default.
gcloud compute networks subnets update default \
--region=${REGION} \
--enable-private-ip-google-access
Vous pouvez vérifier que l'accès privé à Google est activé en exécutant la commande suivante, qui renvoie True ou False.
gcloud compute networks subnets describe default \
--region=${REGION} \
--format="get(privateIpGoogleAccess)"
Créer un bucket de stockage
Créez un bucket de stockage qui sera utilisé pour stocker les éléments créés dans cet atelier de programmation.
Choisissez un nom pour votre bucket. Les noms de bucket doivent être uniques pour tous les utilisateurs.
export BUCKET=<your-bucket-name>
Créez le bucket dans la région où vous prévoyez d'exécuter vos jobs Spark.
gsutil mb -l ${REGION} gs://${BUCKET}
Vous pouvez voir que votre bucket est disponible dans la console Cloud Storage. Vous pouvez également exécuter gsutil ls pour afficher votre bucket.
Créer un serveur d'historique persistant
L'UI Spark fournit un ensemble complet d'outils de débogage et d'insights sur les jobs Spark. Pour afficher l'UI Spark des jobs Dataproc sans serveur terminés, vous devez créer un cluster Dataproc à nœud unique à utiliser comme serveur d'historique persistant.
Attribuez un nom à votre serveur d'historique persistant (PHS, Persistent History Server).
PHS_CLUSTER_NAME=my-phs
Exécutez la commande suivante.
gcloud dataproc clusters create ${PHS_CLUSTER_NAME} \
--region=${REGION} \
--single-node \
--enable-component-gateway \
--properties=spark:spark.history.fs.logDirectory=gs://${BUCKET}/phs/*/spark-job-history
L'UI Spark et le serveur d'historique persistant seront abordés plus en détail plus loin dans cet atelier de programmation.
3. Exécuter des jobs Spark sans serveur avec Dataproc Batches
Dans cet exemple, vous allez utiliser un ensemble de données provenant de l'ensemble de données public Citi Bike Trips de New York (fourni par NYC City Bikes). NYC Citi Bikes est un système de partage de vélos payant disponible à New York. Vous allez effectuer des transformations simples et afficher les ID des 10 stations Citi Bike les plus populaires. Cet exemple utilise également le connecteur spark-bigquery Open Source pour lire et écrire des données de manière fluide entre Spark et BigQuery.
Clonez le dépôt GitHub suivant et accédez (cd) au répertoire contenant le fichier citibike.py.
git clone https://github.com/GoogleCloudPlatform/devrel-demos.git cd devrel-demos/data-analytics/next-2022-workshop/dataproc-serverless
citibike.py
import sys
from pyspark.sql import SparkSession
from pyspark.sql.functions import col
from pyspark.sql.types import BooleanType
if len(sys.argv) == 1:
print("Please provide a GCS bucket name.")
bucket = sys.argv[1]
table = "bigquery-public-data:new_york_citibike.citibike_trips"
spark = SparkSession.builder \
.appName("pyspark-example") \
.config("spark.jars","gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar") \
.getOrCreate()
df = spark.read.format("bigquery").load(table)
top_ten = df.filter(col("start_station_id") \
.isNotNull()) \
.groupBy("start_station_id") \
.count() \
.orderBy("count", ascending=False) \
.limit(10) \
.cache()
top_ten.show()
top_ten.write.option("header", True).csv(f"gs://{bucket}/citibikes_top_ten_start_station_ids")
Envoyez le job à Serverless Spark à l'aide du SDK Cloud, disponible par défaut dans Cloud Shell. Exécutez la commande suivante dans votre shell, qui utilise le Cloud SDK et l'API Dataproc Batches pour envoyer des jobs Spark sans serveur.
gcloud dataproc batches submit pyspark citibike.py \
--batch=citibike-job \
--region=${REGION} \
--deps-bucket=gs://${BUCKET} \
--jars=gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar \
--history-server-cluster=projects/${GOOGLE_CLOUD_PROJECT}/regions/${REGION}/clusters/${PHS_CLUSTER_NAME} \
-- ${BUCKET}
Voici ce que cela implique :
gcloud dataproc batches submitfait référence à l'API Dataproc Batches.pysparkindique que vous envoyez un job PySpark.--batchest le nom du job. Si vous n'en fournissez pas, un UUID aléatoire est généré automatiquement.--region=${REGION}est la région géographique dans laquelle le job sera traité.--deps-bucket=${BUCKET}est l'emplacement où votre fichier Python local est importé avant d'être exécuté dans l'environnement sans serveur.--jars=gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jarinclut le fichier JAR pour le spark-bigquery-connector dans l'environnement d'exécution Spark.--history-server-cluster=projects/${GOOGLE_CLOUD_PROJECT}/regions/${REGION}/clusters/${PHS_CLUSTER}est le nom complet du serveur d'historique persistant. C'est là que les données d'événements Spark (séparées des résultats de la console) sont stockées et peuvent être affichées à partir de l'UI Spark.- Le
--final indique que tout ce qui suit sera des arguments d'exécution pour le programme. Dans ce cas, vous envoyez le nom de votre bucket, comme requis par le job.
Le résultat suivant s'affiche une fois le lot envoyé.
Batch [citibike-job] submitted.
Après quelques minutes, vous obtenez le résultat suivant, ainsi que les métadonnées du job.
+----------------+------+ |start_station_id| count| +----------------+------+ | 519|551078| | 497|423334| | 435|403795| | 426|384116| | 293|372255| | 402|367194| | 285|344546| | 490|330378| | 151|318700| | 477|311403| +----------------+------+ Batch [citibike-job] finished.
Dans la section suivante, vous apprendrez à localiser les journaux de ce job.
Autres fonctionnalités
Avec Spark sans serveur, vous disposez d'options supplémentaires pour exécuter vos jobs.
- Vous pouvez créer une image Docker personnalisée sur laquelle votre job s'exécute. C'est un excellent moyen d'inclure des dépendances supplémentaires, y compris des bibliothèques Python et R.
- Vous pouvez connecter une instance Dataproc Metastore à votre job pour accéder aux métadonnées Hive.
- Si vous avez besoin de plus de contrôle, sachez que Dataproc sans serveur accepte la configuration d'un petit ensemble de propriétés Spark.
4. Métriques et observabilité Dataproc
La console Dataproc Batches liste tous vos jobs Dataproc sans serveur. Dans la console, vous verrez l'ID de lot, l'emplacement, l'état, la date et heure de création, le temps écoulé et le type de chaque job. Cliquez sur l'ID de lot de votre tâche pour en savoir plus.
Cette page comporte des informations telles que Surveillance, qui indique le nombre d'exécuteurs Spark par lot utilisés par votre job au fil du temps (indiquant le niveau d'autoscaling).
Dans l'onglet Détails, vous trouverez d'autres métadonnées sur le job, y compris les arguments et les paramètres qui ont été envoyés avec le job.
Vous pouvez également accéder à tous les journaux depuis cette page. Lorsque des jobs Dataproc sans serveur sont exécutés, trois ensembles de journaux différents sont générés :
- Au niveau du service
- Sortie vers la console
- Journalisation des événements Spark
Au niveau du service : inclut les journaux générés par le service Dataproc sans serveur. Cela inclut, par exemple, les requêtes de Dataproc sans serveur pour des processeurs supplémentaires pour l'autoscaling. Pour les afficher, cliquez sur Afficher les journaux, ce qui ouvrira Cloud Logging.
La sortie de la console est visible sous Sortie.Il s'agit du résultat généré par le job, y compris les métadonnées que Spark imprime au début d'un job ou les instructions d'impression intégrées au job.
La journalisation des événements Spark est accessible depuis l'UI Spark. Comme vous avez fourni un serveur d'historique persistant à votre job Spark, vous pouvez accéder à l'UI Spark en cliquant sur Afficher le serveur d'historique Spark, qui contient des informations sur vos jobs Spark exécutés précédemment. Pour en savoir plus sur l'UI Spark, consultez la documentation Spark officielle.
5. Modèles Dataproc : BQ → GCS
Les modèles Dataproc sont des outils Open Source qui permettent de simplifier davantage les tâches de traitement de données dans le cloud. Ils servent de wrapper pour Dataproc sans serveur et incluent des modèles pour de nombreuses tâches d'importation et d'exportation de données, y compris :
BigQuerytoGCSetGCStoBigQueryGCStoBigTableGCStoJDBCetJDBCtoGCSHivetoBigQueryMongotoGCSetGCStoMongo
La liste complète est disponible dans le fichier README.
Dans cette section, vous allez utiliser les modèles Dataproc pour exporter des données de BigQuery vers GCS.
Cloner le dépôt
Clonez le dépôt et accédez au dossier python.
git clone https://github.com/GoogleCloudPlatform/dataproc-templates.git cd dataproc-templates/python
Configurer l'environnement
Vous allez maintenant définir des variables d'environnement. Les modèles Dataproc utilisent la variable d'environnement GCP_PROJECT pour votre ID de projet. Définissez-la sur GOOGLE_CLOUD_PROJECT..
export GCP_PROJECT=${GOOGLE_CLOUD_PROJECT}
Votre région doit être définie dans l'environnement précédent. Sinon, définissez-le ici.
export REGION=<region>
Les modèles Dataproc utilisent spark-bigquery-connector pour traiter les jobs BigQuery et nécessitent que l'URI soit inclus dans une variable d'environnement JARS. Définissez la variable JARS.
export JARS="gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar"
Configurer les paramètres du modèle
Définissez le nom d'un bucket intermédiaire que le service doit utiliser.
export GCS_STAGING_LOCATION=gs://${BUCKET}
Vous allez ensuite définir des variables spécifiques à la tâche. Pour la table d'entrée, vous ferez à nouveau référence à l'ensemble de données BigQuery NYC Citibike.
BIGQUERY_GCS_INPUT_TABLE=bigquery-public-data.new_york_citibike.citibike_trips
Vous pouvez choisir csv, parquet, avro ou json. Pour cet atelier de programmation, choisissez "CSV". La section suivante explique comment utiliser les modèles Dataproc pour convertir les types de fichiers.
BIGQUERY_GCS_OUTPUT_FORMAT=csv
Définissez le mode de sortie sur overwrite. Vous pouvez choisir entre overwrite, append, ignore ou errorifexists..
BIGQUERY_GCS_OUTPUT_MODE=overwrite
Définissez l'emplacement de sortie GCS sur un chemin d'accès dans votre bucket.
BIGQUERY_GCS_OUTPUT_LOCATION=gs://${BUCKET}/BQtoGCS
Exécuter le modèle
Exécutez le modèle BIGQUERYTOGCS en le spécifiant ci-dessous et en fournissant les paramètres d'entrée que vous avez définis.
./bin/start.sh \
-- --template=BIGQUERYTOGCS \
--bigquery.gcs.input.table=${BIGQUERY_GCS_INPUT_TABLE} \
--bigquery.gcs.output.format=${BIGQUERY_GCS_OUTPUT_FORMAT} \
--bigquery.gcs.output.mode=${BIGQUERY_GCS_OUTPUT_MODE} \
--bigquery.gcs.output.location=${BIGQUERY_GCS_OUTPUT_LOCATION}
Le résultat sera assez bruyant, mais au bout d'une minute environ, vous verrez ce qui suit.
Batch [5766411d6c78444cb5e80f305308d8f8] submitted. ... Batch [5766411d6c78444cb5e80f305308d8f8] finished.
Vous pouvez vérifier que les fichiers ont été générés en exécutant la commande suivante.
gsutil ls ${BIGQUERY_GCS_OUTPUT_LOCATION}
Par défaut, Spark écrit dans plusieurs fichiers, en fonction de la quantité de données. Dans ce cas, environ 30 fichiers générés devraient s'afficher. Les noms des fichiers de sortie Spark sont mis en forme avec le préfixe part, suivi d'un nombre à cinq chiffres (qui indique la référence) et d'une chaîne de hachage. Pour les grandes quantités de données, Spark écrit généralement plusieurs fichiers. Voici un exemple de nom de fichier : part-00000-cbf69737-867d-41cc-8a33-6521a725f7a0-c000.csv.
6. Modèles Dataproc : CSV vers Parquet
Vous allez maintenant utiliser les modèles Dataproc pour convertir des données dans GCS d'un type de fichier à un autre à l'aide de GCSTOGCS. Ce modèle utilise SparkSQL et permet également d'envoyer une requête SparkSQL à traiter lors de la transformation pour un traitement supplémentaire.
Confirmer les variables d'environnement
Vérifiez que GCP_PROJECT, REGION et GCS_STAGING_BUCKET sont définis à partir de la section précédente.
echo ${GCP_PROJECT}
echo ${REGION}
echo ${GCS_STAGING_LOCATION}
Définir des paramètres de modèle
Vous allez maintenant définir les paramètres de configuration pour GCStoGCS. Commencez par l'emplacement des fichiers d'entrée. Notez qu'il s'agit d'un répertoire et non d'un fichier spécifique, car tous les fichiers du répertoire seront traités. Définissez cette valeur sur BIGQUERY_GCS_OUTPUT_LOCATION.
GCS_TO_GCS_INPUT_LOCATION=${BIGQUERY_GCS_OUTPUT_LOCATION}
Définissez le format du fichier d'entrée.
GCS_TO_GCS_INPUT_FORMAT=csv
Définissez le format de sortie souhaité. Vous pouvez choisir entre les formats Parquet, JSON, Avro ou CSV.
GCS_TO_GCS_OUTPUT_FORMAT=parquet
Définissez le mode de sortie sur overwrite. Vous pouvez choisir entre overwrite, append, ignore ou errorifexists..
GCS_TO_GCS_OUTPUT_MODE=overwrite
Définissez l'emplacement de sortie.
GCS_TO_GCS_OUTPUT_LOCATION=gs://${BUCKET}/GCStoGCS
Exécuter le modèle
Exécutez le modèle GCStoGCS.
./bin/start.sh \
-- --template=GCSTOGCS \
--gcs.to.gcs.input.location=${GCS_TO_GCS_INPUT_LOCATION} \
--gcs.to.gcs.input.format=${GCS_TO_GCS_INPUT_FORMAT} \
--gcs.to.gcs.output.format=${GCS_TO_GCS_OUTPUT_FORMAT} \
--gcs.to.gcs.output.mode=${GCS_TO_GCS_OUTPUT_MODE} \
--gcs.to.gcs.output.location=${GCS_TO_GCS_OUTPUT_LOCATION}
Le résultat sera assez bruyant, mais après environ une minute, un message de réussite semblable à celui ci-dessous devrait s'afficher.
Batch [c198787ba8e94abc87e2a0778c05ec8a] submitted. ... Batch [c198787ba8e94abc87e2a0778c05ec8a] finished.
Vous pouvez vérifier que les fichiers ont été générés en exécutant la commande suivante.
gsutil ls ${GCS_TO_GCS_OUTPUT_LOCATION}
Ce modèle vous permet également de fournir des requêtes SparkSQL en transmettant gcs.to.gcs.temp.view.name et gcs.to.gcs.sql.query au modèle. Vous pouvez ainsi exécuter une requête SparkSQL sur les données avant de les écrire dans GCS.
7. Effectuer un nettoyage des ressources
Pour éviter que des frais inutiles ne soient facturés sur votre compte GCP une fois cet atelier de programmation terminé :
- Supprimez le bucket Cloud Storage associé à l'environnement que vous avez créé.
gsutil rm -r gs://${BUCKET}
- Supprimez le cluster Dataproc utilisé pour votre serveur d'historique persistant.
gcloud dataproc clusters delete ${PHS_CLUSTER_NAME} \
--region=${REGION}
- Supprimez les jobs Dataproc sans serveur. Accédez à la console Batches, cochez la case à côté de chaque job que vous souhaitez supprimer, puis cliquez sur SUPPRIMER.
Si vous avez créé un projet spécifiquement pour cet atelier de programmation, vous pouvez également le supprimer :
- Dans la console GCP, accédez à la page Projets.
- Dans la liste des projets, sélectionnez celui que vous souhaitez supprimer, puis cliquez sur "Supprimer".
- Dans la boîte de dialogue, saisissez l'ID du projet, puis cliquez sur "Arrêter" pour supprimer le projet.
8. Étape suivante
Les ressources suivantes vous fournissent d'autres moyens de profiter de Serverless Spark :
- Découvrez comment orchestrer les workflows Dataproc sans serveur à l'aide de Cloud Composer.
- Découvrez comment intégrer Dataproc sans serveur aux pipelines Kubeflow.