Transformer des données et les écrire dans une table Apache Iceberg

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 :

  1. Transférez les données Parquet des courses de taxi à New York.
  2. Écrivez le script de job PySpark.
  3. Envoyez le job par lot sans serveur.
  4. 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.

  1. 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 the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  2. Verify that billing is enabled for your Google Cloud project.

  3. 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.enable permission. 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.

    Enable the APIs

  4. 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.

  5. Créez un bucket Cloud Storage :
    gcloud storage buckets create gs://BUCKET_NAME
    Remplacez BUCKET_NAME par 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 :

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 :

É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.

  1. 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-user
    
  2. Vé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.

  1. Définissez les propriétés du catalogue Iceberg dans un fichier d'indicateurs YAML.

    Créez iceberg-flags.yaml contenant 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
    EOF
    
    Proprié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é.
    type Définit le type de catalogue pour utiliser la spécification du catalogue REST Apache Iceberg (rest).
    uri URL du point de terminaison de l'API REST pour BigLake Metastore.
    io-impl Configure 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-project Transmet 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.
    warehouse URI de ressource BigLake (bl://projects/...) pointant vers la ressource de catalogue spécifique dans votre projet Google Cloud .
    rest.auth.type Authentifie automatiquement les appels d'API de catalogue REST à l'aide d'identifiants Google Cloud (org.apache.iceberg.gcp.auth.GoogleAuthManager).
    spark.sql.extensions Active les extensions Iceberg SQL et la compatibilité avec DataFrameWriterV2 (writeTo) dans Spark SQL.
  2. 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-bucket pour préproduire le fichier de script dans Cloud Storage avant de démarrer le job.
    • Lors de votre premier envoi 3.0 d'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 de DataprocRMExecutorsAllocator dans les journaux du pilote lors de la suppression de l'exécuteur. Ces avertissements sont temporaires et non fatals, à condition que le lot atteigne SUCCEEDED.

É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.

  1. Inspecter les fichiers Iceberg dans Cloud Storage.

    Vérifiez que l'entrepôt Iceberg contient les répertoires metadata/ et data/ attendus :

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. 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}"
    
  3. 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