Ce tutoriel explique comment lire les enregistrements de trajets en taxi à New York à partir d'un fichier Parquet sur Cloud Storage, les transformer avec Apache Spark sur le runtime sans serveur Managed Service pour Apache Spark 3.0 et écrire les résultats agrégés dans une table Apache Iceberg à l'aide du catalogue REST BigLake.
Dans ce tutoriel, vous allez accomplir les tâches suivantes :
- Transférez les données Parquet des courses de taxi à New York.
- Écrivez le script de job PySpark.
- Envoyez le job par lot sans serveur.
- Vérifiez et interrogez la table Iceberg.
Avant de commencer
Configurez votre projet et effectuez d'autres tâches de démarrage.
Configurer votre projet Google Cloud
Configurez votre projet selon vos besoins pour activer les API, attribuer des rôles Identity and Access Management (IAM), authentifier les Identifiants par défaut de l'application et créer un bucket Cloud Storage.
Activer les API
Utilisez la console Google Cloud pour activer les API requises.
-
In the Google Cloud console, on the project selector page, select or create a Google Cloud project.
Roles required to select or create a project
- Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
-
Create a project: To create a project, you need the Project Creator role
(
roles/resourcemanager.projectCreator), which contains theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that billing is enabled for your Google Cloud project.
Enable the Dataproc, Cloud Storage, BigQuery, BigLake, Cloud Logging, Compute Engine, Cloud Resource Manager, and Dataproc Resource Manager APIs.
Roles required to enable APIs
To enable APIs, you need the
serviceusage.services.enablepermission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.-
Créez des identifiants d'authentification locaux pour votre compte utilisateur :
gcloud auth application-default login
Si une erreur d'authentification est renvoyée et que vous utilisez un fournisseur d'identité (IdP) externe, vérifiez que vous vous êtes connecté à la gcloud CLI avec votre identité fédérée.
-
Créez un bucket Cloud Storage :
Remplacezgcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEpar un nom qui répond aux conditions requises pour le nom des buckets.
Accorder des rôles IAM si nécessaire
Certains rôles IAM sont requis pour exécuter les exemples sur cette page. En fonction des règles d'administration, ces rôles peuvent déjà avoir été accordés. Pour vérifier les attributions de rôles, consultez Do you need to grant roles? (Devez-vous attribuer des rôles ?).
Pour en savoir plus sur l'attribution de rôles, consultez Gérer l'accès aux projets, aux dossiers et aux organisations.
Rôles utilisateur
Par défaut, l'environnement d'exécution sans serveur 3.0 de Managed Service pour Apache Spark s'exécute avec vos identifiants d'utilisateur final. Les rôles de compte de service ne sont pas obligatoires. Pour en savoir plus, consultez Personas et rôles IAM sans serveur.
Pour obtenir les autorisations nécessaires pour envoyer une charge de travail par lot sans serveur, demandez à votre administrateur de vous accorder les rôles IAM suivants :
-
Exécutez des charges de travail sur Runtime 3.x (EUC par défaut) :
- Éditeur Dataproc sans serveur (
roles/dataproc.serverlessEditor) sur le projet - Administrateur BigLake (
roles/biglake.admin) sur le projet - Administrateur BigQuery (
roles/bigquery.admin) sur le projet - Administrateur des objets Storage (
roles/storage.objectAdmin) sur le projet
- Éditeur Dataproc sans serveur (
Votre administrateur peut exécuter le script bash suivant pour attribuer des rôles à votre compte utilisateur.
```bash
for ROLE in \
roles/dataproc.serverlessEditor \
roles/biglake.admin \
roles/bigquery.admin \
roles/storage.objectAdmin
do
gcloud projects add-iam-policy-binding "${PROJECT_ID}" \
--member="user:${USER_ACCOUNT}" \
--role="${ROLE}"
done
```
Configurer les variables d'environnement
Exécutez le script Bash suivant pour définir les variables d'environnement de shell utilisées tout au long de ce tutoriel.
# 1. Active project ID, user account, and project number
export PROJECT_ID="$(gcloud config get-value project)"
export USER_ACCOUNT="$(gcloud config get-value account)"
export PROJECT_NUMBER="$(gcloud projects describe "${PROJECT_ID}" --format="value(projectNumber)")"
# 2. Regional deployment and networking settings
export REGION="us-central1"
export SUBNET_NAME="default"
export SUBNET_RANGE="10.128.0.0/20"
# 3. Storage and BigLake Iceberg catalog resources
export BUCKET_NAME="BUCKET_NAME"
export OUTPUT_CATALOG_NAME="lakehouse"
export OUTPUT_DATASET_NAME="nyc_taxi"
export OUTPUT_TABLE_NAME="yellow_trips_analyzed"
export INPUT_PARQUET_PATH="gs://${BUCKET_NAME}/raw/nyc_taxi/yellow_tripdata_2024-01.parquet"
Remplacez les éléments suivants :
BUCKET_NAME: nom du bucket Cloud Storage que vous avez créé dans Configurer votre projet Google Cloud .
Étape 1 : Créer un catalogue BigLake Iceberg
Le catalogue BigLake Iceberg REST permet aux charges de travail Spark et à BigQuery de découvrir, lire et écrire des tables Iceberg.
Créez le catalogue BigLake Iceberg.
Créez le catalogue avec son emplacement par défaut pointant vers le chemin d'accès à votre entrepôt dans Cloud Storage.
gcloud biglake iceberg catalogs create "${OUTPUT_CATALOG_NAME}" \ --project="${PROJECT_ID}" \ --catalog-type=biglake \ --default-location="gs://${BUCKET_NAME}/warehouse" \ --credential-mode=end-userVérifiez l'état du catalogue.
gcloud biglake iceberg catalogs describe "${OUTPUT_CATALOG_NAME}" --project="${PROJECT_ID}"
Étape 2 : Entreposer les données Parquet des taxis new-yorkais
La NYC Taxi & Limousine Commission (TLC) publie des enregistrements mensuels des courses au format Parquet. Téléchargez un mois de données et copiez-les dans votre bucket Cloud Storage.
curl -L -o yellow_tripdata_2024-01.parquet \
https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2024-01.parquet
gcloud storage cp yellow_tripdata_2024-01.parquet "${INPUT_PARQUET_PATH}"
gcloud storage ls "${INPUT_PARQUET_PATH}"
Étape 3 : Écrire le script de job PySpark
Créez un taxi_to_iceberg.py. Ce script lit les enregistrements Parquet bruts, filtre les trajets non valides, calcule les agrégats quotidiens de passagers, de tarifs, de pourboires et de revenus, et écrit une table Iceberg à l'aide du catalogue REST BigLake.
"""Read NYC taxi Parquet from Cloud Storage and write aggregates to an Iceberg table."""
import sys
from pyspark.sql import SparkSession
from pyspark.sql.functions import (
avg,
col,
count,
round as spark_round,
sum as spark_sum,
to_date,
when,
)
def main():
if len(sys.argv) < 5:
print(
"Usage: taxi_to_iceberg.py <input_parquet_path> <catalog_name>"
" <dataset_name> <table_name>"
)
sys.exit(1)
input_path, catalog_name, dataset_name, table_name = sys.argv[1:5]
full_table = f"`{catalog_name}`.{dataset_name}.{table_name}"
# Iceberg extensions are enabled here; catalog properties are passed at
# submit time via --flags-file (see Step 4).
spark = (
SparkSession.builder.appName("NYC Taxi Parquet to Iceberg")
.config(
"spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
)
.getOrCreate()
)
print(f"Reading raw Parquet data from: {input_path}")
raw_df = spark.read.parquet(input_path)
raw_df.printSchema()
cleaned_df = (
raw_df.filter(
(col("trip_distance") > 0)
& (col("fare_amount") > 0)
& (col("passenger_count") > 0)
)
.withColumn("trip_date", to_date(col("tpep_pickup_datetime")))
.withColumn(
"tip_pct",
when(
col("fare_amount") > 0,
spark_round((col("tip_amount") / col("fare_amount")) * 100, 2),
).otherwise(0.0),
)
)
aggregated_df = (
cleaned_df.groupBy("trip_date", "payment_type")
.agg(
count("*").alias("total_trips"),
spark_sum("passenger_count").alias("total_passengers"),
spark_round(avg("trip_distance"), 2).alias("avg_distance"),
spark_round(avg("fare_amount"), 2).alias("avg_fare"),
spark_round(avg("tip_pct"), 2).alias("avg_tip_percentage"),
spark_round(spark_sum("total_amount"), 2).alias("total_revenue"),
)
.orderBy("trip_date", "payment_type")
)
aggregated_df.show(10, truncate=False)
# Create the namespace, then write with DataFrameWriterV2 (writeTo).
spark.sql(f"CREATE NAMESPACE IF NOT EXISTS `{catalog_name}`.{dataset_name}")
spark.catalog.setCurrentCatalog(catalog_name)
aggregated_df.writeTo(full_table).using("iceberg").createOrReplace()
print(f"Wrote Iceberg table: {full_table}")
spark.stop()
if __name__ == "__main__":
main()
Important : Lorsque vous écrivez dans un catalogue BigLake Iceberg, utilisez toujours l'API writeTo (DataFrameWriterV2), définissez d'abord le catalogue actuel, assurez-vous que l'espace de noms existe et entourez le nom du catalogue d'accents graves dans les appels SQL et writeTo.
Étape 4 : Envoyer le job par lot sans serveur
Définissez les propriétés du catalogue dans un fichier de configuration et envoyez le job par lot à Managed Service pour Apache Spark sans serveur.
Définissez les propriétés du catalogue Iceberg dans un fichier d'indicateurs YAML.
Créez
iceberg-flags.yamlcontenant les propriétés du catalogue REST BigLake. La Google Cloud CLI lit ce fichier à l'aide de l'argument--flags-file.cat <<EOF > iceberg-flags.yaml --properties: spark.sql.catalog.${OUTPUT_CATALOG_NAME}: org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.${OUTPUT_CATALOG_NAME}.type: rest spark.sql.catalog.${OUTPUT_CATALOG_NAME}.uri: https://biglake.googleapis.com/iceberg/v1/restcatalog spark.sql.catalog.${OUTPUT_CATALOG_NAME}.io-impl: org.apache.iceberg.gcp.gcs.GCSFileIO spark.sql.catalog.${OUTPUT_CATALOG_NAME}.header.x-goog-user-project: ${PROJECT_ID} spark.sql.catalog.${OUTPUT_CATALOG_NAME}.warehouse: bl://projects/${PROJECT_ID}/catalogs/${OUTPUT_CATALOG_NAME} spark.sql.catalog.${OUTPUT_CATALOG_NAME}.rest.auth.type: org.apache.iceberg.gcp.auth.GoogleAuthManager spark.sql.extensions: org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions EOFPropriété Objectif / Description spark.sql.catalog.${OUTPUT_CATALOG_NAME}Enregistre le plug-in de catalogue Apache Iceberg ( org.apache.iceberg.spark.SparkCatalog) pour le nom de catalogue personnalisé.typeDéfinit le type de catalogue pour utiliser la spécification du catalogue REST Apache Iceberg ( rest).uriURL du point de terminaison de l'API REST pour BigLake Metastore. io-implConfigure Iceberg pour utiliser Cloud Storage FileIO ( org.apache.iceberg.gcp.gcs.GCSFileIO) pour les opérations de lecture/écriture de données et de métadonnées hautes performances.header.x-goog-user-projectTransmet l'ID de votre projet Google Cloud en tant qu'en-tête de requête pour l'attribution du quota et de la facturation de l'API BigLake. warehouseURI de ressource BigLake ( bl://projects/...) pointant vers la ressource de catalogue spécifique dans votre projet Google Cloud .rest.auth.typeAuthentifie automatiquement les appels d'API de catalogue REST à l'aide d'identifiants Google Cloud ( org.apache.iceberg.gcp.auth.GoogleAuthManager).spark.sql.extensionsActive les extensions Iceberg SQL et la compatibilité avec DataFrameWriterV2(writeTo) dans Spark SQL.Envoyez le job par lot PySpark.
Définissez les variables de version d'exécution et de table cible, puis envoyez le job par lot à l'aide de
gcloud dataproc batches submit pyspark:export RUNTIME_VERSION="3.0" export SCRIPT_FILE="taxi_to_iceberg.py" gcloud dataproc batches submit pyspark "${SCRIPT_FILE}" \ --flags-file=iceberg-flags.yaml \ --project="${PROJECT_ID}" \ --region="${REGION}" \ --version="${RUNTIME_VERSION}" \ --subnet="${SUBNET_NAME}" \ --deps-bucket="gs://${BUCKET_NAME}" \ -- \ "${INPUT_PARQUET_PATH}" \ "${OUTPUT_CATALOG_NAME}" \ "${OUTPUT_DATASET_NAME}" \ "${OUTPUT_TABLE_NAME}"Remarques :
- Lorsque vous envoyez un script Python local, gcloud CLI utilise
--deps-bucketpour préproduire le fichier de script dans Cloud Storage avant de démarrer le job. - Lors de votre premier envoi
3.0d'environnement d'exécution à l'aide de vos identifiants d'utilisateur final, vous recevrez une invite d'autorisation OAuth unique. Accordez l'accès et renvoyez le fichier. - Lorsque vous utilisez le runtime
3.0, vous pouvez observer des avertissements deDataprocRMExecutorsAllocatordans les journaux du pilote lors de la suppression de l'exécuteur. Ces avertissements sont temporaires et non fatals, à condition que le lot atteigneSUCCEEDED.
- Lorsque vous envoyez un script Python local, gcloud CLI utilise
Étape 5 : Vérifier et interroger la table Iceberg
Vérifiez les fichiers de sortie dans Cloud Storage, confirmez l'enregistrement de la table dans BigLake et interrogez la table à l'aide de BigQuery.
Inspecter les fichiers Iceberg dans Cloud Storage.
Vérifiez que l'entrepôt Iceberg contient les répertoires
metadata/etdata/attendus :gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"Vérifiez l'enregistrement de la table auprès de BigLake.
gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \ --catalog="${OUTPUT_CATALOG_NAME}" \ --namespace="${OUTPUT_DATASET_NAME}" \ --project="${PROJECT_ID}"Interrogez BigQuery.
BigLake enregistre automatiquement la table Iceberg dans BigQuery :
bq query \ --project_id="${PROJECT_ID}" \ --location="${REGION}" \ --use_legacy_sql=false \ "SELECT trip_date, payment_type, total_trips, avg_fare, avg_tip_percentage, total_revenue FROM \`${PROJECT_ID}.${OUTPUT_CATALOG_NAME}.${OUTPUT_DATASET_NAME}.${OUTPUT_TABLE_NAME}\` ORDER BY total_trips DESC LIMIT 10"
Effectuer un nettoyage
Pour éviter que des frais ne soient facturés sur votre compte Google Cloud , supprimez les ressources que vous avez créées dans ce tutoriel.
# 1. Delete BigLake Iceberg catalog, and table metadata
gcloud biglake iceberg catalogs delete "${OUTPUT_CATALOG_NAME}" \
--project="${PROJECT_ID}" --quiet 2>/dev/null || true
# 2. Delete BigQuery dataset (if created)
bq rm -r -f -d "${PROJECT_ID}:${OUTPUT_DATASET_NAME}" 2>/dev/null || true
# 3. Delete Cloud Storage bucket and warehouse files
gcloud storage rm -r "gs://${BUCKET_NAME}"
Étapes suivantes
- En savoir plus sur les charges de travail par lot sans serveur Managed Service pour Apache Spark
- Découvrez les personas et les rôles IAM sans serveur.
- En savoir plus sur les tables BigLake Iceberg dans BigQuery
- Consultez Résoudre les problèmes courants liés au Lakehouse.