Daten transformieren und in eine Apache Iceberg-Tabelle schreiben

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:

  1. NYC-Taxi-Parquet-Daten bereitstellen
  2. Schreiben Sie das PySpark-Jobskript.
  3. Senden Sie den serverlosen Batch-Job.
  4. 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.

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

  5. Erstellen Sie einen Cloud Storage-Bucket:
    gcloud storage buckets create gs://BUCKET_NAME
    Ersetzen Sie BUCKET_NAME durch 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:

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:

Schritt 1: BigLake Iceberg-Katalog erstellen

Mit dem BigLake Iceberg REST-Katalog können Spark-Arbeitslasten und BigQuery Iceberg-Tabellen erkennen, lesen und schreiben.

  1. 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-user
    
  2. Katalogstatus 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.

  1. Definieren Sie Iceberg-Katalogattribute in einer YAML-Flag-Datei.

    Erstellen Sie iceberg-flags.yaml mit 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
    EOF
    
    Attribut 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.
    type Legt den Katalogtyp für die Verwendung der Apache Iceberg REST-Katalogspezifikation (rest) fest.
    uri Die REST API-Endpunkt-URL für BigLake REST Metastore.
    io-impl Konfiguriert 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.
    warehouse Der BigLake-Ressourcen-URI (bl://projects/...), der auf die spezifische Katalogressource in Ihrem Google Cloud -Projekt verweist.
    rest.auth.type Authentifiziert REST-Katalog-API-Aufrufe automatisch mit Google Cloud -Anmeldedaten (org.apache.iceberg.gcp.auth.GoogleAuthManager).
    spark.sql.extensions Aktiviert Iceberg-SQL-Erweiterungen und die Unterstützung von DataFrameWriterV2 (writeTo) in Spark SQL.
  2. 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.0 verwenden, werden möglicherweise Warnungen von DataprocRMExecutorsAllocator in den Treiberlogs während des Executor-Teardown angezeigt. Diese Warnungen sind vorübergehend und nicht schwerwiegend, sofern der Batch SUCCEEDED erreicht.

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.

  1. Iceberg-Dateien in Cloud Storage prüfen

    Prüfen Sie, ob das Iceberg-Warehouse die erwarteten Verzeichnisse metadata/ und data/ enthält:

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. Tabellenregistrierung mit BigLake prüfen

    gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \
      --catalog="${OUTPUT_CATALOG_NAME}" \
      --namespace="${OUTPUT_DATASET_NAME}" \
      --project="${PROJECT_ID}"
    
  3. 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}"

Nächste Schritte