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:
- Organiza los datos de Parquet de los taxis de la ciudad de Nueva York.
- Escribe la secuencia de comandos del trabajo de PySpark.
- Envía el trabajo por lotes sin servidores.
- 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.
-
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.-
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.
-
Crea un bucket de Cloud Storage:
Reemplazagcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEpor 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:
-
Ejecuta cargas de trabajo en el entorno de ejecución 3.x (EUC predeterminado):
- Editor de Dataproc Serverless (
roles/dataproc.serverlessEditor) en el proyecto - Administrador de BigLake (
roles/biglake.admin) en el proyecto - Administrador de BigQuery (
roles/bigquery.admin) en el proyecto - Administrador de objetos de almacenamiento (
roles/storage.objectAdmin) en el proyecto
- Editor de Dataproc Serverless (
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:
BUCKET_NAME: Es el nombre del bucket de Cloud Storage que creaste en Configura tu proyecto de Google Cloud .
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.
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-userVerifica 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.
Define las propiedades del catálogo de Iceberg en un archivo de marcas YAML.
Crea
iceberg-flags.yamlque 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 EOFPropiedad 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.typeEstablece el tipo de catálogo para usar la especificación del catálogo de REST de Apache Iceberg ( rest).uriEs la URL del extremo de API de REST de BigLake Metastore. io-implConfigura 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-projectPasa 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. warehouseEl URI del recurso de BigLake ( bl://projects/...) que apunta al recurso de catálogo específico en tu proyecto Google Cloud .rest.auth.typeAutentica 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.extensionsHabilita las extensiones de Iceberg SQL y la compatibilidad con DataFrameWriterV2(writeTo) en Spark SQL.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-bucketpara 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.0del 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.0en el tiempo de ejecución, es posible que observes advertencias deDataprocRMExecutorsAllocatoren los registros del controlador durante el cierre del ejecutor. Estas advertencias son transitorias y no fatales, siempre que el lote llegue aSUCCEEDED.
- Cuando envías una secuencia de comandos de Python local, la gcloud CLI usa
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.
Inspecciona los archivos de Iceberg en Cloud Storage.
Verifica que el almacén de Iceberg contenga los directorios
metadata/ydata/esperados:gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"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}"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?
- Obtén más información sobre las cargas de trabajo por lotes sin servidores de Managed Service para Apache Spark.
- Explora los arquetipos y los roles de IAM sin servidores.
- Obtén más información sobre las tablas de BigLake Iceberg en BigQuery.
- Consulta Cómo solucionar problemas habituales de Lakehouse.