Este tutorial mostra como ler registros de viagens de táxi da cidade de Nova York (NYC) de um arquivo Parquet no Cloud Storage, transformá-los com o Apache Spark no ambiente de execução sem servidor do Serviço Gerenciado para Apache Spark 3.0 e gravar os resultados agregados em uma tabela do Apache Iceberg usando o catálogo REST do BigLake.
Neste tutorial, você irá:
- Faça o staging dos dados do Parquet de táxis de Nova York.
- Escreva o script do job do PySpark.
- Envie o job em lote sem servidor.
- Verifique e consulte a tabela do Iceberg.
Antes de começar
Configure seu projeto e realize outras tarefas de inicialização.
Configurar o projeto do Google Cloud
Configure o projeto conforme necessário para ativar APIs, conceder papéis do Identity and Access Management (IAM), autenticar Application Default Credentials e criar um bucket do Cloud Storage.
Ativar APIs
Use o console Google Cloud para ativar as APIs necessárias.
-
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.-
Crie credenciais de autenticação local para sua conta de usuário:
gcloud auth application-default login
Se um erro de autenticação for retornado e você estiver usando um provedor de identidade (IdP) externo, confirme se você fez login na CLI gcloud com sua identidade federada.
-
Crie um bucket do Cloud Storage:
Substituagcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEpor um nome de bucket que atenda aos requisitos de nomenclatura de bucket.
Conceder papéis do IAM, se necessário
Alguns papéis do IAM são necessários para executar os exemplos nesta página. Dependendo das políticas da organização, essas funções já podem ter sido concedidas. Para verificar as concessões de papéis, consulte Você precisa conceder papéis?.
Para mais informações sobre a concessão de papéis, consulte Gerenciar o acesso a projetos, pastas e organizações.
Funções do usuário
Por padrão, o ambiente de execução sem servidor do Serviço Gerenciado para Apache Spark 3.0 é executado com suas credenciais de usuário final (EUC). Não é necessário ter papéis de conta de serviço. Para
mais informações, consulte Personas e papéis do IAM sem servidor.
Para receber as permissões necessárias para enviar uma carga de trabalho em lote sem servidor, peça ao administrador para conceder a você os seguintes papéis do IAM:
-
Execute cargas de trabalho no ambiente de execução 3.x (EUC padrão):
- Editor do Dataproc sem servidor (
roles/dataproc.serverlessEditor) no projeto - Administrador do BigLake (
roles/biglake.admin) no projeto - Administrador do BigQuery (
roles/bigquery.admin) no projeto - Administrador de objetos do Storage (
roles/storage.objectAdmin) no projeto
- Editor do Dataproc sem servidor (
O administrador pode executar o seguinte script bash para conceder papéis à sua conta de usuário.
```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 as variáveis de ambiente
Execute o script bash a seguir para definir as variáveis ambiente shell usadas neste 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"
Substitua:
BUCKET_NAME: o nome do bucket do Cloud Storage que você criou em Configurar seu projeto Google Cloud .
Etapa 1. Criar um catálogo do BigLake Iceberg
O catálogo REST do BigLake Iceberg permite que cargas de trabalho do Spark e o BigQuery descubram, leiam e gravem tabelas do Iceberg.
Crie o catálogo do BigLake Iceberg.
Crie o catálogo com o local padrão apontando para o caminho do seu data warehouse no 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-userVerifique a integridade do catálogo.
gcloud biglake iceberg catalogs describe "${OUTPUT_CATALOG_NAME}" --project="${PROJECT_ID}"
Etapa 2. Preparar os dados Parquet de táxis de Nova York
A Comissão de Táxis e Limusines de Nova York (TLC, na sigla em inglês) publica registros mensais de viagens em Parquet. Faça o download de um mês de dados e copie para o bucket do 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}"
Etapa 3. Escrever o script do job do PySpark
Criar taxi_to_iceberg.py. Esse script lê os registros brutos do Parquet, filtra viagens inválidas, calcula os agregados diários de passageiros, tarifas, gorjetas e receita e grava uma tabela do Iceberg usando o catálogo REST do 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:ao gravar em um catálogo do BigLake Iceberg, sempre use a API writeTo (DataFrameWriterV2), defina o catálogo atual primeiro, verifique se o namespace existe e coloque o nome do catálogo entre crases em chamadas SQL e writeTo.
Etapa 4. Enviar o job em lote sem servidor
Defina as propriedades do catálogo em um arquivo de configuração e envie o job em lote para o Serviço Gerenciado para Apache Spark sem servidor.
Defina as propriedades do catálogo do Iceberg em um arquivo de flags YAML.
Crie
iceberg-flags.yamlque contenha as propriedades do catálogo REST do BigLake. A CLI do Google Cloud lê esse arquivo usando o 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 EOFPropriedade Finalidade / Descrição spark.sql.catalog.${OUTPUT_CATALOG_NAME}Registra o plug-in de catálogo do Apache Iceberg ( org.apache.iceberg.spark.SparkCatalog) para o nome do catálogo personalizado.typeDefine o tipo de catálogo para usar a especificação do catálogo REST do Apache Iceberg ( rest).uriO URL do endpoint de API REST para o BigLake REST Metastore. io-implConfigura o Iceberg para usar o Cloud Storage FileIO ( org.apache.iceberg.gcp.gcs.GCSFileIO) em operações de leitura/gravação de dados e metadados de alta performance.header.x-goog-user-projectTransmite o ID do seu projeto Google Cloud como um cabeçalho de solicitação para atribuição de cota e faturamento da API BigLake. warehouseO URI do recurso do BigLake ( bl://projects/...) que aponta para o recurso de catálogo específico no seu projeto Google Cloud .rest.auth.typeAutentica automaticamente as chamadas da API REST do catálogo usando credenciais Google Cloud ( org.apache.iceberg.gcp.auth.GoogleAuthManager).spark.sql.extensionsAtiva as extensões SQL do Iceberg e o suporte a DataFrameWriterV2(writeTo) no Spark SQL.Envie o job em lote do PySpark.
Defina as variáveis de versão do ambiente de execução e da tabela de destino e envie o job em lote usando
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}"Observações:
- Ao enviar um script Python local, a CLI gcloud usa
--deps-bucketpara preparar o arquivo de script no Cloud Storage antes de iniciar o job. - Na primeira vez que você enviar uma
3.0usando suas credenciais de usuário final (EUC), vai receber uma solicitação de consentimento do OAuth única. Conceda acesso e envie novamente. - Ao usar o
3.0de execução, talvez você veja avisos doDataprocRMExecutorsAllocatornos registros de driver durante o encerramento do executor. Esses avisos são temporários e não fatais, supondo que o lote alcanceSUCCEEDED.
- Ao enviar um script Python local, a CLI gcloud usa
Etapa 5: verificar e consultar a tabela Iceberg
Verifique os arquivos de saída no Cloud Storage, confirme o registro da tabela no BigLake e consulte a tabela usando o BigQuery.
Inspecione os arquivos do Iceberg no Cloud Storage.
Verifique se o data warehouse do Iceberg contém os diretórios
metadata/edata/esperados:gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"Verifique o registro da tabela com o BigLake.
gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \ --catalog="${OUTPUT_CATALOG_NAME}" \ --namespace="${OUTPUT_DATASET_NAME}" \ --project="${PROJECT_ID}"Consulta do BigQuery.
O BigLake registra automaticamente a tabela do Iceberg no 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"
Limpar
Para evitar cobranças na sua conta do Google Cloud , exclua os recursos criados neste 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}"
A seguir
- Saiba mais sobre as cargas de trabalho em lote sem servidor do Serviço Gerenciado para Apache Spark.
- Conheça Personas e papéis do IAM sem servidor.
- Leia sobre as tabelas do BigLake Iceberg no BigQuery.
- Consulte Solução de problemas comuns do Lakehouse.