In dieser Anleitung wird gezeigt, wie Sie Taxifahrten in New York City (NYC) aus einer Parquet-Datei in Cloud Storage lesen, sie mit Apache Spark in der serverlosen Laufzeit 3.0 von Managed Service for Apache Spark transformieren und die aggregierten Ergebnisse mit dem BigLake REST-Katalog in eine Apache Iceberg-Tabelle schreiben.
In dieser Anleitung führen Sie die folgenden Aufgaben aus:
- NYC-Taxi-Parquet-Daten bereitstellen
- Schreiben Sie das PySpark-Jobskript.
- Senden Sie den serverlosen Batch-Job.
- Iceberg-Tabelle prüfen und abfragen
Hinweis
Richten Sie Ihr Projekt ein und führen Sie andere Startaufgaben aus.
Projekt in Google Cloud einrichten
Richten Sie Ihr Projekt nach Bedarf ein, um APIs zu aktivieren, IAM-Rollen (Identity and Access Management) zu gewähren, Standardanmeldedaten für Anwendungen zu authentifizieren und einen Cloud Storage-Bucket zu erstellen.
APIs aktivieren
Aktivieren Sie die erforderlichen APIs in der Google Cloud Console.
-
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.-
Erstellen Sie lokale Anmeldedaten zur Authentifizierung für Ihr Nutzerkonto:
gcloud auth application-default login
Wenn ein Authentifizierungsfehler zurückgegeben wird und Sie einen externen Identitätsanbieter (IdP) verwenden, prüfen Sie, ob Sie sich mit Ihrer föderierten Identität in der gcloud CLI angemeldet haben.
-
Erstellen Sie einen Cloud Storage-Bucket:
Ersetzen Siegcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEdurch einen Bucket-Namen, der den Anforderungen für Bucket-Namen entspricht.
Bei Bedarf IAM-Rollen zuweisen
Für die Ausführung der Beispiele auf dieser Seite sind bestimmte IAM-Rollen erforderlich. Je nach Organisationsrichtlinien wurden diese Rollen möglicherweise bereits gewährt. Informationen zum Prüfen von Rollenzuweisungen finden Sie unter Müssen Sie Rollen zuweisen?.
Weitere Informationen zum Zuweisen von Rollen finden Sie unter Zugriff auf Projekte, Ordner und Organisationen verwalten.
Nutzerrollen
Standardmäßig wird die serverlose Laufzeit von Managed Service for Apache Spark 3.0 mit Ihren Endnutzeranmeldedaten ausgeführt. Dienstkontorollen sind nicht erforderlich. Weitere Informationen finden Sie unter Personas und serverlose IAM-Rollen.
Bitten Sie Ihren Administrator, Ihnen die folgenden IAM-Rollen zuzuweisen, damit Sie die nötigen Berechtigungen zum Einreichen einer serverlosen Batcharbeitslast haben:
-
Arbeitslasten in Runtime 3.x (Standard-EUC) ausführen:
- Dataproc Serverless Editor (
roles/dataproc.serverlessEditor) für das Projekt - BigLake-Administrator (
roles/biglake.admin) für das Projekt - BigQuery Admin (
roles/bigquery.admin) für das Projekt - Storage-Objekt-Administrator (
roles/storage.objectAdmin) für das Projekt
- Dataproc Serverless Editor (
Ihr Administrator kann das folgende Bash-Skript ausführen, um Ihrem Nutzerkonto Rollen zuzuweisen.
```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
```
Umgebungsvariablen konfigurieren
Führen Sie das folgende Bash-Skript aus, um Shell-Umgebungsvariablen festzulegen, die in dieser Anleitung verwendet werden.
# 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"
Ersetzen Sie Folgendes:
BUCKET_NAME: Der Name des Cloud Storage-Bucket, den Sie unter Google Cloud -Projekt einrichten erstellt haben.
Schritt 1: BigLake Iceberg-Katalog erstellen
Mit dem BigLake Iceberg REST-Katalog können Spark-Arbeitslasten und BigQuery Iceberg-Tabellen erkennen, lesen und schreiben.
Erstellen Sie den BigLake Iceberg-Katalog.
Erstellen Sie den Katalog mit dem standardmäßigen Standort, der auf Ihren Warehouse-Pfad in Cloud Storage verweist.
gcloud biglake iceberg catalogs create "${OUTPUT_CATALOG_NAME}" \ --project="${PROJECT_ID}" \ --catalog-type=biglake \ --default-location="gs://${BUCKET_NAME}/warehouse" \ --credential-mode=end-userKatalogstatus prüfen
gcloud biglake iceberg catalogs describe "${OUTPUT_CATALOG_NAME}" --project="${PROJECT_ID}"
Schritt 2: NYC-Taxidaten im Parquet-Format bereitstellen
Die NYC Taxi & Limousine Commission (TLC) veröffentlicht monatliche Fahrtenaufzeichnungen im Parquet-Format. Laden Sie die Daten für einen Monat herunter und kopieren Sie sie in Ihren Cloud Storage-Bucket.
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}"
Schritt 3: PySpark-Job-Script schreiben
Erstellen Sie taxi_to_iceberg.py. Dieses Skript liest die rohen Parquet-Datensätze, filtert ungültige Fahrten heraus, berechnet tägliche Aggregate für Fahrgäste, Fahrpreis, Trinkgeld und Umsatz und schreibt eine Iceberg-Tabelle mit dem BigLake REST-Katalog.
"""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()
Wichtig:Wenn Sie in einen BigLake Iceberg-Katalog schreiben, verwenden Sie immer die writeTo API (DataFrameWriterV2), legen Sie zuerst den aktuellen Katalog fest, sorgen Sie dafür, dass der Namespace vorhanden ist, und setzen Sie den Katalognamen in SQL- und writeTo-Aufrufen in Backticks.
Schritt 4: Serverlosen Batch-Job senden
Definieren Sie die Katalogeigenschaften in einer Konfigurationsdatei und senden Sie den Batchjob an Managed Service for Apache Spark Serverless.
Definieren Sie Iceberg-Katalogattribute in einer YAML-Flag-Datei.
Erstellen Sie
iceberg-flags.yamlmit den BigLake REST-Katalogattributen. Die Google Cloud CLI liest diese Datei mit dem 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 EOFAttribut Zweck / Beschreibung spark.sql.catalog.${OUTPUT_CATALOG_NAME}Registriert das Katalog-Plug-in von Apache Iceberg ( org.apache.iceberg.spark.SparkCatalog) für den benutzerdefinierten Katalognamen.typeLegt den Katalogtyp für die Verwendung der Apache Iceberg REST-Katalogspezifikation ( rest) fest.uriDie REST API-Endpunkt-URL für BigLake REST Metastore. io-implKonfiguriert Iceberg für die Verwendung von Cloud Storage FileIO ( org.apache.iceberg.gcp.gcs.GCSFileIO) für leistungsstarke Lese-/Schreibvorgänge für Daten und Metadaten.header.x-goog-user-projectÜbergibt Ihre Google Cloud Projekt-ID als Anfrageheader für die BigLake API-Kontingent- und Abrechnungszuordnung. warehouseDer BigLake-Ressourcen-URI ( bl://projects/...), der auf die spezifische Katalogressource in Ihrem Google Cloud -Projekt verweist.rest.auth.typeAuthentifiziert REST-Katalog-API-Aufrufe automatisch mit Google Cloud -Anmeldedaten ( org.apache.iceberg.gcp.auth.GoogleAuthManager).spark.sql.extensionsAktiviert Iceberg-SQL-Erweiterungen und die Unterstützung von DataFrameWriterV2(writeTo) in Spark SQL.Senden Sie den PySpark-Batchjob.
Legen Sie die Variablen für die Laufzeitversion und die Zieltabelle fest und senden Sie den Batchjob dann mit
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}"Hinweise:
- Wenn Sie ein lokales Python-Skript einreichen, verwendet die gcloud CLI
--deps-bucket, um die Skriptdatei in Cloud Storage bereitzustellen, bevor der Job gestartet wird. - Bei der ersten
3.0-Einreichung zur Laufzeit mit Ihren Endnutzeranmeldedaten erhalten Sie eine einmalige OAuth-Zustimmungsaufforderung. Gewähren Sie den Zugriff und reichen Sie den Antrag noch einmal ein. - Wenn Sie die Laufzeit
3.0verwenden, werden möglicherweise Warnungen vonDataprocRMExecutorsAllocatorin den Treiberlogs während des Executor-Teardown angezeigt. Diese Warnungen sind vorübergehend und nicht schwerwiegend, sofern der BatchSUCCEEDEDerreicht.
- Wenn Sie ein lokales Python-Skript einreichen, verwendet die gcloud CLI
Schritt 5: Iceberg-Tabelle überprüfen und abfragen
Prüfen Sie die Ausgabedateien in Cloud Storage, bestätigen Sie die Tabellenregistrierung in BigLake und fragen Sie die Tabelle mit BigQuery ab.
Iceberg-Dateien in Cloud Storage prüfen
Prüfen Sie, ob das Iceberg-Warehouse die erwarteten Verzeichnisse
metadata/unddata/enthält:gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"Tabellenregistrierung mit BigLake prüfen
gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \ --catalog="${OUTPUT_CATALOG_NAME}" \ --namespace="${OUTPUT_DATASET_NAME}" \ --project="${PROJECT_ID}"Abfrage aus BigQuery
BigLake registriert die Iceberg-Tabelle automatisch 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"
Bereinigen
Löschen Sie die in dieser Anleitung erstellten Ressourcen, um zu vermeiden, dass Ihrem Google Cloud Konto Gebühren in Rechnung gestellt werden.
# 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}"