Transforma datos y escribe en una tabla de Apache Iceberg

En este instructivo, se muestra cómo leer registros de viajes en taxi de la ciudad de Nueva York (NYC) desde un archivo Parquet en Cloud Storage, transformarlos con Apache Spark en el entorno de ejecución sin servidores de Managed Service para Apache Spark 3.0 y escribir los resultados agregados en una tabla de Apache Iceberg con el catálogo de REST de BigLake.

En este tutorial, completarás las siguientes tareas:

  1. Organiza los datos de Parquet de los taxis de la ciudad de Nueva York.
  2. Escribe la secuencia de comandos del trabajo de PySpark.
  3. Envía el trabajo por lotes sin servidores.
  4. Verifica y consulta la tabla de Iceberg.

Antes de comenzar

Configura tu proyecto y realiza otras tareas de inicio.

Configura tu proyecto de Google Cloud

Configura tu proyecto según sea necesario para habilitar APIs, otorgar roles de Identity and Access Management (IAM), autenticar credenciales predeterminadas de la aplicación y crear un bucket de Cloud Storage.

Habilita las APIs

Usa la consola de Google Cloud para habilitar las APIs requeridas.

  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 credenciales de autenticación locales para tu cuenta de usuario:

    gcloud auth application-default login

    Si se devuelve un error de autenticación y usas un proveedor de identidad (IdP) externo, confirma que accediste a la gcloud CLI con tu identidad federada.

  5. Crea un bucket de Cloud Storage:
    gcloud storage buckets create gs://BUCKET_NAME
    Reemplaza BUCKET_NAME por un nombre de bucket que cumpla con los requisitos de nombres de buckets.

Otorga roles de IAM si es necesario

Se requieren ciertos roles de IAM para ejecutar los ejemplos de esta página. Según las políticas de la organización, es posible que estos roles ya se hayan otorgado. Para verificar las asignaciones de roles, consulta ¿Necesitas otorgar roles?.

Para obtener más información sobre cómo otorgar roles, consulta Administra el acceso a proyectos, carpetas y organizaciones.

Roles de usuario

De forma predeterminada, el entorno de ejecución sin servidores 3.0 de Managed Service para Apache Spark se ejecuta con tus credenciales de usuario final (EUC). No se requieren roles de cuenta de servicio. Para obtener más información, consulta Arquetipos y roles de IAM sin servidores.

Para obtener los permisos que necesitas para enviar una carga de trabajo por lotes sin servidores, pídele a tu administrador que te otorgue los siguientes roles de IAM:

Tu administrador puede ejecutar la siguiente secuencia de comandos de Bash para otorgar roles a tu cuenta de usuario.

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

Configure las variables de entorno

Ejecuta la siguiente secuencia de comandos en Bash para establecer las variables de entorno de shell que se usan en este instructivo.

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

Reemplaza lo siguiente:

Paso 1: Crea un catálogo de BigLake Iceberg

El catálogo de REST de BigLake Iceberg permite que las cargas de trabajo de Spark y BigQuery descubran, lean y escriban tablas de Iceberg.

  1. Crea el catálogo de BigLake Iceberg.

    Crea el catálogo con su ubicación predeterminada que apunta a la ruta de acceso de tu almacén en 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 el estado del catálogo.

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

Paso 2: Almacena en etapa intermedia los datos de Parquet de los taxis de la ciudad de Nueva York

La Comisión de Taxis y Limusinas (TLC) de la ciudad de Nueva York publica registros de viajes mensuales en formato Parquet. Descarga los datos de un mes y cópialos en tu bucket de 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}"

Paso 3: Escribe la secuencia de comandos del trabajo de PySpark

Crea taxi_to_iceberg.py. Esta secuencia de comandos lee los registros sin procesar de Parquet, filtra los viajes no válidos, calcula los agregados diarios de pasajeros, tarifas, propinas y ganancias, y escribe una tabla de Iceberg con el catálogo REST de 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: Cuando escribas en un catálogo de BigLake Iceberg, siempre usa la API de writeTo (DataFrameWriterV2), primero establece el catálogo actual, asegúrate de que exista el espacio de nombres y encierra el nombre del catálogo entre comillas inversas en las llamadas de SQL y writeTo.

Paso 4: Envía el trabajo por lotes sin servidores

Define las propiedades del catálogo en un archivo de configuración y envía el trabajo por lotes a Managed Service para Apache Spark sin servidores.

  1. Define las propiedades del catálogo de Iceberg en un archivo de marcas YAML.

    Crea iceberg-flags.yaml que contenga las propiedades del catálogo de la API de BigLake REST. Google Cloud CLI lee este archivo con el argumento --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
    
    Propiedad Propósito o descripción
    spark.sql.catalog.${OUTPUT_CATALOG_NAME} Registra el complemento de catálogo de Apache Iceberg (org.apache.iceberg.spark.SparkCatalog) para el nombre del catálogo personalizado.
    type Establece el tipo de catálogo para usar la especificación del catálogo de REST de Apache Iceberg (rest).
    uri Es la URL del extremo de API de REST de BigLake Metastore.
    io-impl Configura Iceberg para que use Cloud Storage FileIO (org.apache.iceberg.gcp.gcs.GCSFileIO) para operaciones de lectura y escritura de metadatos y datos de alto rendimiento.
    header.x-goog-user-project Pasa tu ID del proyecto Google Cloud como encabezado de la solicitud para la atribución de facturación y cuota de la API de BigLake.
    warehouse El URI del recurso de BigLake (bl://projects/...) que apunta al recurso de catálogo específico en tu proyecto Google Cloud .
    rest.auth.type Autentica automáticamente las llamadas a la API de REST del catálogo con credenciales de Google Cloud (org.apache.iceberg.gcp.auth.GoogleAuthManager).
    spark.sql.extensions Habilita las extensiones de Iceberg SQL y la compatibilidad con DataFrameWriterV2 (writeTo) en Spark SQL.
  2. Envía el trabajo por lotes de PySpark.

    Configura las variables de versión del entorno de ejecución y de tabla de destino, y, luego, envía el trabajo por lotes con 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}"
    

    Notas:

    • Cuando envías una secuencia de comandos de Python local, la gcloud CLI usa --deps-bucket para almacenar en etapa intermedia el archivo de secuencia de comandos en Cloud Storage antes de iniciar el trabajo.
    • En el primer envío de 3.0 del entorno de ejecución con tus credenciales de usuario final (EUC), recibirás una solicitud de consentimiento de OAuth única. Otorga acceso y vuelve a enviar la solicitud.
    • Cuando usas 3.0 en el tiempo de ejecución, es posible que observes advertencias de DataprocRMExecutorsAllocator en los registros del controlador durante el cierre del ejecutor. Estas advertencias son transitorias y no fatales, siempre que el lote llegue a SUCCEEDED.

Paso 5: Verifica y consulta la tabla de Iceberg

Verifica los archivos de salida en Cloud Storage, confirma el registro de la tabla en BigLake y consulta la tabla con BigQuery.

  1. Inspecciona los archivos de Iceberg en Cloud Storage.

    Verifica que el almacén de Iceberg contenga los directorios metadata/ y data/ esperados:

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. Verifica el registro de la tabla con BigLake.

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

    BigLake registra automáticamente la tabla de Iceberg en 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"
    

Realiza una limpieza

Para evitar que se generen cargos en tu cuenta de Google Cloud , borra los recursos que creaste en este instructivo.

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

¿Qué sigue?