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à:
- Organizza i dati Parquet dei taxi di New York.
- Scrivi lo script del job PySpark.
- Invia il job batch serverless.
- 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.
-
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 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.
-
Crea un bucket Cloud Storage:
Sostituiscigcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEcon 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:
-
Esegui carichi di lavoro su Runtime 3.x (EUC predefinito):
- Dataproc Serverless Editor (
roles/dataproc.serverlessEditor) sul progetto - BigLake Admin (
roles/biglake.admin) sul progetto - Amministratore BigQuery (
roles/bigquery.admin) sul progetto - Storage Object Admin (
roles/storage.objectAdmin) sul progetto
- Dataproc Serverless Editor (
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:
BUCKET_NAME: il nome del bucket Cloud Storage che hai creato in Configura il progetto Google Cloud .
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.
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-userVerifica 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.
Definisci le proprietà del catalogo Iceberg in un file di flag YAML.
Crea
iceberg-flags.yamlche 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 EOFProprietà 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.typeImposta il tipo di catalogo in modo da utilizzare la specifica del catalogo REST Apache Iceberg ( rest).uriL'URL dell'endpoint API REST per BigLake REST Metastore. io-implConfigura 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-projectTrasmette l'ID progetto Google Cloud come intestazione della richiesta per l'attribuzione di quota e fatturazione dell'API BigLake. warehouseL'URI della risorsa BigLake ( bl://projects/...) che punta alla risorsa catalogo specifica nel tuo progetto Google Cloud .rest.auth.typeAutentica automaticamente le chiamate all'API REST Catalog utilizzando le credenziali Google Cloud ( org.apache.iceberg.gcp.auth.GoogleAuthManager).spark.sql.extensionsAttiva le estensioni SQL Iceberg e il supporto di DataFrameWriterV2(writeTo) in Spark SQL.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-bucketper preparare il file di script in Cloud Storage prima di avviare il job. - Al primo invio di
3.0utilizzando le credenziali dell'utente finale (EUC), riceverai una richiesta di consenso OAuth una tantum. Concedi l'accesso e invia di nuovo. - Quando utilizzi
3.0in fase di runtime, potresti notare avvisi daDataprocRMExecutorsAllocatornei log dei driver durante l'arresto anomalo dell'executor. Questi avvisi sono temporanei e non fatali, supponendo che il batch raggiungaSUCCEEDED.
- Quando invia uno script Python locale, gcloud CLI utilizza
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.
Ispeziona i file Iceberg in Cloud Storage.
Verifica che il warehouse Iceberg contenga le directory
metadata/edata/previste:gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"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}"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
- Scopri di più sui workload batch serverless di Managed Service for Apache Spark.
- Esplora le identità e i ruoli IAM serverless.
- Scopri di più sulle tabelle BigLake Iceberg in BigQuery.
- Consulta la sezione Risoluzione dei problemi comuni di Lakehouse.