Trasforma i dati e scrivili in una tabella Apache Iceberg

Questo tutorial mostra come leggere i record delle corse dei taxi di New York City (NYC) da un file Parquet su Cloud Storage, trasformarli con Apache Spark nel runtime serverless di Managed Service for Apache Spark 3.0 e scrivere i risultati aggregati in una tabella Apache Iceberg utilizzando il catalogo BigLake REST.

In questo tutorial completerai le seguenti attività:

  1. Organizza i dati Parquet dei taxi di New York.
  2. Scrivi lo script del job PySpark.
  3. Invia il job batch serverless.
  4. Verifica ed esegui query sulla tabella Iceberg.

Prima di iniziare

Configura il progetto ed esegui altre attività di avvio.

Configura il progetto Google Cloud

Configura il progetto in base alle esigenze per abilitare le API, concedere i ruoli Identity and Access Management (IAM), autenticare le Credenziali predefinite dell'applicazione e creare un bucket Cloud Storage.

Abilita API

Utilizza la Google Cloud console per abilitare le API richieste.

  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. Crea credenziali di autenticazione locali per il tuo account utente:

    gcloud auth application-default login

    Se viene restituito un errore di autenticazione e utilizzi un provider di identità (IdP) esterno, verifica di aver acceduto a gcloud CLI con la tua identità federata.

  5. Crea un bucket Cloud Storage:
    gcloud storage buckets create gs://BUCKET_NAME
    Sostituisci BUCKET_NAME con un nome del bucket che soddisfi i requisiti per la denominazione dei bucket.

Concedi ruoli IAM, se necessario

Per eseguire gli esempi in questa pagina sono necessari determinati ruoli IAM. A seconda delle norme dell'organizzazione, questi ruoli potrebbero essere già stati concessi. Per controllare le concessioni dei ruoli, consulta Devi concedere ruoli?

Per saperne di più sulla concessione dei ruoli, consulta Gestisci l'accesso a progetti, cartelle e organizzazioni.

Ruoli utente

Per impostazione predefinita, il runtime serverless di Managed Service for Apache Spark 3.0 viene eseguito con le credenziali dell'utente finale. I ruoli del service account non sono obbligatori. Per saperne di più, consulta Personas e ruoli IAM serverless.

Per ottenere le autorizzazioni necessarie per inviare un carico di lavoro batch serverless, chiedi all'amministratore di concederti i seguenti ruoli IAM:

L'amministratore può eseguire il seguente script bash per concedere ruoli al tuo account utente.

```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
```

Configura le variabili di ambiente

Esegui questo script bash per impostare le variabili di ambiente della shell utilizzate in questo tutorial.

# 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"

Sostituisci quanto segue:

Passaggio 1: Crea un catalogo BigLake Iceberg

Il catalogo REST BigLake Iceberg consente ai carichi di lavoro Spark e BigQuery di rilevare, leggere e scrivere tabelle Iceberg.

  1. Crea il catalogo BigLake Iceberg.

    Crea il catalogo con la posizione predefinita che punta al percorso del warehouse in 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. Verifica l'integrità del catalogo.

    gcloud biglake iceberg catalogs describe "${OUTPUT_CATALOG_NAME}" --project="${PROJECT_ID}"
    

Passaggio 2: Organizza i dati Parquet dei taxi di New York

La NYC Taxi & Limousine Commission (TLC) pubblica mensilmente le registrazioni delle corse in Parquet. Scarica un mese di dati e copiali nel 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}"

Passaggio 3: Scrivi lo script del job PySpark

Crea taxi_to_iceberg.py. Questo script legge i record Parquet non elaborati, filtra i viaggi non validi, calcola gli aggregati giornalieri di passeggeri, tariffe, mance e entrate e scrive una tabella Iceberg utilizzando il catalogo 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()

Importante:quando scrivi in un catalogo BigLake Iceberg, utilizza sempre l'API writeTo (DataFrameWriterV2), imposta prima il catalogo corrente, assicurati che lo spazio dei nomi esista e racchiudi il nome del catalogo tra apici inversi nelle chiamate SQL e writeTo.

Passaggio 4: Invia il job batch serverless

Definisci le proprietà del catalogo in un file di configurazione e invia il job batch a Managed Service for Apache Spark serverless.

  1. Definisci le proprietà del catalogo Iceberg in un file di flag YAML.

    Crea iceberg-flags.yaml che contenga le proprietà del catalogo BigLake REST. Google Cloud CLI legge questo file utilizzando l'argomento --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
    
    Proprietà Scopo / Descrizione
    spark.sql.catalog.${OUTPUT_CATALOG_NAME} Registra il plug-in del catalogo di Apache Iceberg (org.apache.iceberg.spark.SparkCatalog) per il nome del catalogo personalizzato.
    type Imposta il tipo di catalogo in modo da utilizzare la specifica del catalogo REST Apache Iceberg (rest).
    uri L'URL dell'endpoint API REST per BigLake REST Metastore.
    io-impl Configura Iceberg per utilizzare Cloud Storage FileIO (org.apache.iceberg.gcp.gcs.GCSFileIO) per operazioni di lettura/scrittura di dati e metadati ad alte prestazioni.
    header.x-goog-user-project Trasmette l'ID progetto Google Cloud come intestazione della richiesta per l'attribuzione di quota e fatturazione dell'API BigLake.
    warehouse L'URI della risorsa BigLake (bl://projects/...) che punta alla risorsa catalogo specifica nel tuo progetto Google Cloud .
    rest.auth.type Autentica automaticamente le chiamate all'API REST Catalog utilizzando le credenziali Google Cloud (org.apache.iceberg.gcp.auth.GoogleAuthManager).
    spark.sql.extensions Attiva le estensioni SQL Iceberg e il supporto di DataFrameWriterV2 (writeTo) in Spark SQL.
  2. Invia il job batch PySpark.

    Imposta le variabili della versione runtime e della tabella di destinazione, poi invia il job batch utilizzando 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}"
    

    Note:

    • Quando invia uno script Python locale, gcloud CLI utilizza --deps-bucket per preparare il file di script in Cloud Storage prima di avviare il job.
    • Al primo invio di 3.0 utilizzando le credenziali dell'utente finale (EUC), riceverai una richiesta di consenso OAuth una tantum. Concedi l'accesso e invia di nuovo.
    • Quando utilizzi 3.0 in fase di runtime, potresti notare avvisi da DataprocRMExecutorsAllocator nei log dei driver durante l'arresto anomalo dell'executor. Questi avvisi sono temporanei e non fatali, supponendo che il batch raggiunga SUCCEEDED.

Passaggio 5: verifica ed esegui query sulla tabella Iceberg

Verifica i file di output in Cloud Storage, conferma la registrazione della tabella in BigLake ed esegui query sulla tabella utilizzando BigQuery.

  1. Ispeziona i file Iceberg in Cloud Storage.

    Verifica che il warehouse Iceberg contenga le directory metadata/ e data/ previste:

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. Verifica la registrazione della tabella con BigLake.

    gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \
      --catalog="${OUTPUT_CATALOG_NAME}" \
      --namespace="${OUTPUT_DATASET_NAME}" \
      --project="${PROJECT_ID}"
    
  3. Query da BigQuery.

    BigLake registra automaticamente la tabella Iceberg in 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"
    

Esegui la pulizia

Per evitare che al tuo account Google Cloud vengano addebitati costi, elimina le risorse che hai creato in questo tutorial.

# 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}"

Passaggi successivi