데이터 변환 및 Apache Iceberg 테이블에 쓰기

이 튜토리얼에서는 Cloud Storage의 Parquet 파일에서 뉴욕시 (NYC) 택시 이동 기록을 읽고, Apache Spark용 관리형 서비스 서버리스 런타임 3.0에서 Apache Spark로 변환하고, BigLake REST 카탈로그를 사용하여 집계된 결과를 Apache Iceberg 테이블에 쓰는 방법을 보여줍니다.

이 튜토리얼에서는 다음 태스크를 완료합니다.

  1. NYC 택시 Parquet 데이터를 스테이징합니다.
  2. PySpark 작업 스크립트를 작성합니다.
  3. 서버리스 일괄 작업을 제출합니다.
  4. Iceberg 테이블을 확인하고 쿼리합니다.

시작하기 전에

프로젝트를 설정하고 기타 시작 작업을 실행합니다.

Google Cloud 프로젝트 설정

API를 사용 설정하고, Identity and Access Management (IAM) 역할을 부여하고, 애플리케이션 기본 사용자 인증 정보를 인증하고, Cloud Storage 버킷을 만드는 데 필요한 대로 프로젝트를 설정합니다.

API 사용 설정

Google Cloud 콘솔을 사용하여 필요한 API를 사용 설정합니다.

  1. Google Cloud 계정에 로그인합니다. Google Cloud를 처음 사용하는 경우 계정을 만들고 Google 제품의 실제 성능을 평가해 보세요. 신규 고객에게는 워크로드를 실행, 테스트, 배포하는 데 사용할 수 있는 $300의 무료 크레딧이 제공됩니다.
  2. 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

  3. Verify that billing is enabled for your Google Cloud project.

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

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

  6. Verify that billing is enabled for your Google Cloud project.

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

  8. 로컬 셸을 사용하는 경우 사용자 계정에 대한 로컬 인증 사용자 인증 정보를 만듭니다.

    gcloud auth application-default login

    Cloud Shell을 사용하는 경우 이 작업을 수행할 필요는 없습니다.

    인증 오류가 반환되고 외부 ID 공급업체(IdP)를 사용하는 경우 제휴 ID로 gcloud CLI에 로그인했는지 확인합니다.

  9. Cloud Storage 버킷을 만듭니다.
    gcloud storage buckets create gs://BUCKET_NAME
    BUCKET_NAME버킷 이름 요구사항을 충족하는 버킷 이름으로 바꿉니다.

필요한 경우 IAM 역할 부여

이 페이지의 예시를 실행하려면 특정 IAM 역할이 필요합니다. 조직 정책에 따라 이러한 역할이 이미 부여되었을 수 있습니다. 역할 부여를 확인하려면 역할을 부여해야 하나요?를 참고하세요.

역할 부여에 대한 자세한 내용은 프로젝트, 폴더, 조직에 대한 액세스 관리를 참고하세요.

사용자 역할

기본적으로 Managed Service for Apache Spark 서버리스 런타임 3.0은 최종 사용자 사용자 인증 정보 (EUC)로 실행됩니다. 서비스 계정 역할은 필요하지 않습니다. 자세한 내용은 페르소나 및 서버리스 IAM 역할을 참고하세요.

서버리스 일괄 워크로드를 제출하는 데 필요한 권한을 얻으려면 관리자에게 다음 IAM 역할을 부여해 달라고 요청하세요.

관리자는 다음 bash 스크립트를 실행하여 사용자 계정에 역할을 부여할 수 있습니다.

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

환경 변수 구성

다음 bash 스크립트를 실행하여 이 튜토리얼 전체에서 사용되는 셸 환경 변수를 설정합니다.

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

다음을 바꿉니다.

1단계: BigLake Iceberg 카탈로그 만들기

BigLake Iceberg REST 카탈로그를 사용하면 Spark 워크로드와 BigQuery가 Iceberg 테이블을 검색, 읽기, 쓸 수 있습니다.

  1. BigLake Iceberg 카탈로그를 만듭니다.

    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. 카탈로그 상태를 확인합니다.

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

2단계: NYC 택시 Parquet 데이터 스테이징

뉴욕시 택시 및 리무진 조합 (TLC)은 Parquet 형식으로 월별 운행 기록을 게시합니다. 한 달 분량의 데이터를 다운로드하여 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}"

3단계: PySpark 작업 스크립트 작성

taxi_to_iceberg.py를 만듭니다. 이 스크립트는 원시 Parquet 레코드를 읽고, 유효하지 않은 이동을 필터링하고, 일일 승객, 요금, 팁, 수익 집계를 계산하고, BigLake REST 카탈로그를 사용하여 Iceberg 테이블을 작성합니다.

"""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()

중요: BigLake Iceberg 카탈로그에 쓸 때는 항상 writeTo (DataFrameWriterV2) API를 사용하고, 먼저 현재 카탈로그를 설정하고, 네임스페이스가 있는지 확인하고, SQL 및 writeTo 호출에서 카탈로그 이름을 백틱으로 묶습니다.

4단계: 서버리스 일괄 작업 제출

구성 파일에서 카탈로그 속성을 정의하고 Managed Service for Apache Spark 서버리스에 일괄 작업을 제출합니다.

  1. YAML 플래그 파일에서 Iceberg 카탈로그 속성을 정의합니다.

    BigLake REST 카탈로그 속성이 포함된 iceberg-flags.yaml를 만듭니다. Google Cloud CLI는 --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
    
    속성 목적 / 설명
    spark.sql.catalog.${OUTPUT_CATALOG_NAME} 맞춤 카탈로그 이름에 Apache Iceberg의 카탈로그 플러그인 (org.apache.iceberg.spark.SparkCatalog)을 등록합니다.
    type Apache Iceberg REST 카탈로그 사양 (rest)을 사용하도록 카탈로그 유형을 설정합니다.
    uri BigLake REST Metastore의 REST API 엔드포인트 URL입니다.
    io-impl 고성능 데이터 및 메타데이터 읽기/쓰기 작업을 위해 Cloud Storage FileIO (org.apache.iceberg.gcp.gcs.GCSFileIO)를 사용하도록 Iceberg를 구성합니다.
    header.x-goog-user-project BigLake API 할당량 및 청구 기여 분석을 위해 Google Cloud 프로젝트 ID를 요청 헤더로 전달합니다.
    warehouse Google Cloud 프로젝트의 특정 카탈로그 리소스를 가리키는 BigLake 리소스 URI (bl://projects/...)입니다.
    rest.auth.type Google Cloud 사용자 인증 정보 (org.apache.iceberg.gcp.auth.GoogleAuthManager)를 사용하여 REST 카탈로그 API 호출을 자동으로 인증합니다.
    spark.sql.extensions Spark SQL에서 Iceberg SQL 확장 프로그램과 DataFrameWriterV2 (writeTo) 지원을 사용 설정합니다.
  2. PySpark 일괄 작업을 제출합니다.

    런타임 버전과 타겟 테이블 변수를 설정한 후 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}"
    

    참고:

    • 로컬 Python 스크립트를 제출할 때 gcloud CLI는 --deps-bucket를 사용하여 작업을 시작하기 전에 스크립트 파일을 Cloud Storage에 스테이징합니다.
    • 최종 사용자 인증 정보 (EUC)를 사용하여 첫 번째 런타임 3.0 제출 시 일회성 OAuth 동의 메시지가 표시됩니다. 액세스 권한을 부여하고 다시 제출하세요.
    • 런타임 3.0를 사용하면 실행기 해체 중에 드라이버 로그에 DataprocRMExecutorsAllocator의 경고가 표시될 수 있습니다. 배치가 SUCCEEDED에 도달한다고 가정하면 이러한 경고는 일시적이며 심각하지 않습니다.

5단계: Iceberg 테이블 확인 및 쿼리

Cloud Storage에서 출력 파일을 확인하고, BigLake에서 테이블 등록을 확인하고, BigQuery를 사용하여 테이블을 쿼리합니다.

  1. Cloud Storage에서 Iceberg 파일을 검사합니다.

    Iceberg 웨어하우스에 예상되는 metadata/data/ 디렉터리가 포함되어 있는지 확인합니다.

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. BigLake에 테이블 등록을 확인합니다.

    gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \
      --catalog="${OUTPUT_CATALOG_NAME}" \
      --namespace="${OUTPUT_DATASET_NAME}" \
      --project="${PROJECT_ID}"
    
  3. BigQuery에서 쿼리합니다.

    BigLake는 BigQuery에 Iceberg 테이블을 자동으로 등록합니다.

    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"
    

삭제

Google Cloud 계정에 비용이 청구되지 않게 하려면 이 튜토리얼에서 만든 리소스를 삭제합니다.

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

다음 단계