A configuração do catálogo personalizado do Apache Iceberg para BigQuery para se conectar ao catálogo de ambientes de execução do Lakehouse permite o uso de clusters do Serviço Gerenciado para Apache Spark ou Serviço Gerenciado para Apache Spark para criar um único catálogo compartilhado.
Como parte do Lakehouse sem fronteiras, essa configuração permite usar um cluster do Serviço Gerenciado para Apache Spark ou o Serviço Gerenciado para Apache Spark para criar um único catálogo compartilhado.
Depois de configurada, essa camada de metadados funciona com mecanismos de código aberto, como o Apache Spark e o Apache Flink, para gerenciar seus formatos de tabela aberta.
Antes de começar
- Ative o faturamento no seu Google Cloud projeto. Saiba como verificar se o faturamento está ativado em um projeto.
Ative as APIs do BigQuery e do Serviço Gerenciado para Apache Spark.
Entenda como o catálogo de ambientes de execução do Lakehouse funciona.
Funções exigidas
Para receber as permissões necessárias para configurar o catálogo de ambientes de execução do Lakehouse, peça ao administrador para conceder a você os papéis do IAM a seguir:
-
Criar um cluster do Serviço Gerenciado para Apache Spark:
Worker do Dataproc (
roles/dataproc.worker) na conta de serviço padrão do Compute Engine no projeto -
Criar tabelas do catálogo de ambientes de execução do Lakehouse:
- Worker do Dataproc (
roles/dataproc.worker) na conta de serviço da VM do Serviço Gerenciado para Apache Spark no projeto - Editor de dados do BigQuery (
roles/bigquery.dataEditor) na conta de serviço da VM do Dataproc no projeto - Usuário de objetos do Storage (
roles/storage.objectUser) na conta de serviço da VM do Dataproc no projeto
- Worker do Dataproc (
-
Consultar tabelas do catálogo de ambientes de execução do Lakehouse:
- Visualizador de dados do BigQuery (
roles/bigquery.dataViewer) no projeto - Usuário do BigQuery (
roles/bigquery.user) no projeto - Leitor de objetos do Storage (
roles/storage.objectViewer) no projeto
- Visualizador de dados do BigQuery (
Para mais informações sobre a concessão de papéis, consulte Gerenciar o acesso a projetos, pastas e organizações.
Também é possível conseguir as permissões necessárias com papéis personalizados ou outros papéis predefinidos.
Configurar o catálogo de ambientes de execução do Lakehouse com o Serviço Gerenciado para Apache Spark
É possível configurar o catálogo de ambientes de execução do Lakehouse com o Serviço Gerenciado para Apache Spark usando o Apache Spark ou o Apache Flink:
Apache Spark
Configure um novo cluster. Para criar um cluster do Serviço Gerenciado para Apache Spark, execute o comando
gcloud dataproc clusters createa seguir, que contém as configurações necessárias para usar o catálogo de ambientes de execução do Lakehouse:gcloud dataproc clusters create CLUSTER_NAME \ --project=PROJECT_ID \ --region=LOCATION \ --single-node
Substitua:
CLUSTER_NAME: um nome para o cluster do Serviço Gerenciado para Apache Spark.PROJECT_ID: o ID do Google Cloud projeto em que você está criando o cluster.LOCATION: a região do Compute Engine em que você está criando o cluster.
Envie um job do Apache Spark usando um destes métodos:
Google Cloud CLI
gcloud dataproc jobs submit spark-sql \ --project=PROJECT_ID \ --cluster=CLUSTER_NAME \ --region=REGION \ --jars=https://storage-download.googleapis.com/maven-central/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/1.6.1/iceberg-spark-runtime-3.5_2.12-1.6.1.jar,BIGLAKE_ICEBERG_CATALOG_JAR \ --properties=spark.sql.catalog.CATALOG_NAME=org.apache.iceberg.spark.SparkCatalog,\ spark.sql.catalog.CATALOG_NAME.catalog-impl=org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog,\ spark.sql.catalog.CATALOG_NAME.gcp_project=PROJECT_ID,\ spark.sql.catalog.CATALOG_NAME.gcp_location=LOCATION,\ spark.sql.catalog.CATALOG_NAME.warehouse=WAREHOUSE_DIRECTORY \ --execute="SPARK_SQL_COMMAND"
Substitua:
PROJECT_ID: o ID do Google Cloud projeto que contém o cluster do Serviço Gerenciado para Apache Spark.CLUSTER_NAME: o nome do cluster do Serviço Gerenciado para Apache Spark que você está usando para executar o job do Apache Spark SQL.REGION: a região do Compute Engine em que o cluster está localizado.BIGLAKE_ICEBERG_CATALOG_JAR: o URI do Cloud Storage do plug-in do catálogo personalizado do Apache Iceberg a ser usado. Dependendo do número da versão do Apache Iceberg, selecione uma das seguintes opções:- Apache Iceberg 1.9.1:
gs://spark-lib/bigquery/iceberg-bigquery-catalog-1.9.1-1.0.1.jar - Apache Iceberg 1.6.1:
gs://spark-lib/bigquery/iceberg-bigquery-catalog-1.6.1-1.0.2.jar
- Apache Iceberg 1.9.1:
LOCATION: o local dos recursos do BigQuery.CATALOG_NAME: o nome do catálogo do Apache Spark a ser usado com o job SQL.WAREHOUSE_DIRECTORY: a pasta do Cloud Storage que contém seu data warehouse. Esse valor começa comgs://.SPARK_SQL_COMMAND: a consulta SQL do Apache Spark que você quer executar. Essa consulta inclui os comandos para criar seus recursos. Por exemplo, para criar um namespace e uma tabela.
CLI do spark-sql
No Google Cloud console, acesse a página Instâncias de VM.
Para se conectar a uma instância de VM do Serviço Gerenciado para Apache Spark, clique em SSH na linha que lista o nome da instância de VM principal do cluster do Serviço Gerenciado para Apache Spark, que é o nome do cluster seguido por um sufixo
-m. O resultado será o seguinte:Connected, host fingerprint: ssh-rsa ... Linux cluster-1-m 3.16.0-0.bpo.4-amd64 ... ... example-cluster@cluster-1-m:~$No terminal, execute o seguinte comando de inicialização para o catálogo de ambientes de execução do Lakehouse:
spark-sql \ --jars https://storage-download.googleapis.com/maven-central/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/1.6.1/iceberg-spark-runtime-3.5_2.12-1.6.1.jar,BIGLAKE_ICEBERG_CATALOG_JAR \ --conf spark.sql.catalog.CATALOG_NAME=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.CATALOG_NAME.catalog-impl=org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog \ --conf spark.sql.catalog.CATALOG_NAME.gcp_project=PROJECT_ID \ --conf spark.sql.catalog.CATALOG_NAME.gcp_location=LOCATION \ --conf spark.sql.catalog.CATALOG_NAME.warehouse=WAREHOUSE_DIRECTORY
Substitua:
BIGLAKE_ICEBERG_CATALOG_JAR: o URI do Cloud Storage do plug-in do catálogo personalizado do Apache Iceberg a ser usado. Dependendo do número da versão do Apache Iceberg, selecione uma das seguintes opções:- Apache Iceberg 1.9.1:
gs://spark-lib/bigquery/iceberg-bigquery-catalog-1.9.1-1.0.1.jar - Apache Iceberg 1.6.1:
gs://spark-lib/bigquery/iceberg-bigquery-catalog-1.6.1-1.0.2.jar
- Apache Iceberg 1.9.1:
CATALOG_NAME: o nome do catálogo do Apache Spark que você está usando com o job SQL.PROJECT_ID: o Google Cloud ID do projeto do catálogo de ambientes de execução do Lakehouse ao qual o catálogo do Apache Spark está vinculado.LOCATION: o Google Cloud local do catálogo de ambientes de execução do Lakehouse.WAREHOUSE_DIRECTORY: a pasta do Cloud Storage que contém seu data warehouse. Esse valor começa comgs://.
Depois de se conectar ao cluster, o terminal do Apache Spark mostra o prompt
spark-sql, que pode ser usado para enviar jobs do Apache Spark.spark-sql (default)>
Apache Flink
Crie um cluster do Serviço Gerenciado para Apache Spark com o componente opcional do Apache Flink ativado e verifique se você está usando o Serviço Gerenciado para Apache Spark
2.2ou mais recente.No Google Cloud console do, acesse a página Instâncias de VM.
Na lista de instâncias de máquina virtual, clique em SSH para se conectar à instância de VM principal do cluster do Serviço Gerenciado para Apache Spark, que é listada como o nome do cluster seguido por um sufixo
-m.Configure o plug-in do catálogo personalizado do Apache Iceberg para BigQuery para o catálogo de ambientes de execução do Lakehouse:
FLINK_VERSION=1.19 ICEBERG_VERSION=1.6.1 cd /usr/lib/flink sudo wget -c https://repo.maven.apache.org/maven2/org/apache/iceberg/iceberg-flink-runtime-${FLINK_VERSION}/${ICEBERG_VERSION}/iceberg-flink-runtime-${FLINK_VERSION}-${ICEBERG_VERSION}.jar -P lib sudo gcloud storage cp gs://spark-lib/bigquery/iceberg-bigquery-catalog-${ICEBERG_VERSION}-1.0.2.jar lib/
Inicie a sessão do Apache Flink no YARN:
HADOOP_CLASSPATH=`hadoop classpath` sudo bin/yarn-session.sh -nm flink-dataproc -d sudo bin/sql-client.sh embedded \ -s yarn-session
Crie um catálogo no Apache Flink:
CREATE CATALOG CATALOG_NAME WITH ( 'type'='iceberg', 'warehouse'='WAREHOUSE_DIRECTORY', 'catalog-impl'='org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog', 'gcp_project'='PROJECT_ID', 'gcp_location'='LOCATION' );
Substitua:
CATALOG_NAME: o identificador do catálogo do Apache Flink, que está vinculado ao catálogo de ambientes de execução do Lakehouse.WAREHOUSE_DIRECTORY: o caminho base para o diretório do data warehouse (a pasta do Cloud Storage em que o Apache Flink cria arquivos). Esse valor começa comgs://.PROJECT_ID: o ID do projeto do catálogo de ambientes de execução do Lakehouse ao qual o catálogo do Apache Flink está vinculado.LOCATION: o local dos recursos do BigQuery.
A sessão do Apache Flink agora está conectada ao catálogo de ambientes de execução do Lakehouse, e você pode executar comandos SQL do Apache Flink.
Gerenciar recursos do catálogo de ambientes de execução do Lakehouse
Agora que você está conectado ao catálogo de ambientes de execução do Lakehouse, é possível criar e visualizar recursos com base nos metadados armazenados no catálogo de ambientes de execução do Lakehouse.
Por exemplo, tente executar os comandos a seguir na sessão SQL interativa do Apache Flink para criar um banco de dados e uma tabela do Apache Iceberg.
Use o catálogo personalizado do Apache Iceberg para BigQuery:
USE CATALOG CATALOG_NAME;
Substitua
CATALOG_NAMEpelo identificador do catálogo do Apache Flink.Crie um banco de dados, que cria um conjunto de dados no BigQuery:
CREATE DATABASE IF NOT EXISTS DATABASE_NAME;
Substitua
DATABASE_NAMEpelo nome do novo banco de dados.Use o banco de dados que você criou:
USE DATABASE_NAME;
Crie uma tabela do Apache Iceberg. O comando a seguir cria uma tabela de vendas de exemplo:
CREATE TABLE IF NOT EXISTS ICEBERG_TABLE_NAME ( order_number BIGINT, price DECIMAL(32,2), buyer ROW<first_name STRING, last_name STRING>, order_time TIMESTAMP(3) );
Substitua
ICEBERG_TABLE_NAMEpor um nome para a nova tabela.Ver metadados da tabela:
DESCRIBE EXTENDED ICEBERG_TABLE_NAME;
Listar tabelas no banco de dados:
SHOW TABLES;
Ingerir dados na tabela
Depois de criar uma tabela do Apache Iceberg na seção anterior, você pode usar o Apache Flink DataGen como uma fonte de dados para ingerir dados em tempo real na tabela. As etapas a seguir são um exemplo desse fluxo de trabalho:
Crie uma tabela temporária usando o DataGen:
CREATE TEMPORARY TABLE DATABASE_NAME.TEMP_TABLE_NAME WITH ( 'connector' = 'datagen', 'rows-per-second' = '10', 'fields.order_number.kind' = 'sequence', 'fields.order_number.start' = '1', 'fields.order_number.end' = '1000000', 'fields.price.min' = '0', 'fields.price.max' = '10000', 'fields.buyer.first_name.length' = '10', 'fields.buyer.last_name.length' = '10' ) LIKE DATABASE_NAME.ICEBERG_TABLE_NAME (EXCLUDING ALL);
Substitua:
DATABASE_NAME: o nome do banco de dados para armazenar a tabela temporária.TEMP_TABLE_NAME: um nome para a tabela temporária.ICEBERG_TABLE_NAME: o nome da tabela do Apache Iceberg que você criou na seção anterior.
Defina o paralelismo como 1:
SET 'parallelism.default' = '1';
Defina o intervalo de checkpoint:
SET 'execution.checkpointing.interval' = '10second';
Defina o checkpoint:
SET 'state.checkpoints.dir' = 'hdfs:///flink/checkpoints';
Inicie o job de streaming em tempo real:
INSERT INTO ICEBERG_TABLE_NAME SELECT * FROM TEMP_TABLE_NAME;
O resultado será o seguinte:
[INFO] Submitting SQL update statement to the cluster... [INFO] SQL update statement has been successfully submitted to the cluster: Job ID: 0de23327237ad8a811d37748acd9c10b
Para verificar o status do job de streaming, faça o seguinte:
No Google Cloud console, acesse a página Clusters.
Selecione o cluster.
Clique na guia Interfaces da Web.
Clique no link YARN ResourceManager.
Na interface do YARN ResourceManager , encontre sua sessão do Apache Flink e clique no link ApplicationMaster em Tracking UI.
Na coluna Status, confirme se o status do job é Em execução.
Consulte dados de streaming no cliente SQL do Apache Flink:
SELECT * FROM ICEBERG_TABLE_NAME /*+ OPTIONS('streaming'='true', 'monitor-interval'='3s')*/ ORDER BY order_time desc LIMIT 20;
Consulte dados de streaming no BigQuery:
SELECT * FROM `DATABASE_NAME.ICEBERG_TABLE_NAME` ORDER BY order_time desc LIMIT 20;
Encerre o job de streaming no cliente SQL do Apache Flink:
STOP JOB 'JOB_ID';
Substitua
JOB_IDpelo ID do job que foi exibido na saída quando você criou o job de streaming.
Configurar o catálogo de ambientes de execução do Lakehouse com o Serviço Gerenciado para Apache Spark
É possível configurar o catálogo de ambientes de execução do Lakehouse com o Serviço Gerenciado para Apache Spark usando o Apache Spark SQL ou o PySpark.
SQL do Apache Spark
Crie um arquivo SQL com os comandos SQL do Apache Spark que você quer executar no catálogo de ambientes de execução do Lakehouse. Por exemplo, esse comando cria um namespace e uma tabela:
CREATE NAMESPACE `CATALOG_NAME`.NAMESPACE_NAME; CREATE TABLE `CATALOG_NAME`.NAMESPACE_NAME.TABLE_NAME (id int, data string) USING ICEBERG LOCATION 'WAREHOUSE_DIRECTORY';
Substitua:
CATALOG_NAME: o nome do catálogo que faz referência à tabela do Apache Spark.NAMESPACE_NAME: o nome do namespace que faz referência à tabela do Apache Spark.TABLE_NAME: um nome para a tabela do Apache Spark.WAREHOUSE_DIRECTORY: o URI da pasta do Cloud Storage em que o data warehouse está armazenado.
Envie um job em lote do Apache Spark SQL executando o comando
gcloud dataproc batches submit spark-sqla seguir:gcloud dataproc batches submit spark-sql SQL_SCRIPT_PATH \ --project=PROJECT_ID \ --region=REGION \ --subnet=projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_NAME \ --deps-bucket=BUCKET_PATH \ --properties="spark.sql.catalog.CATALOG_NAME=org.apache.iceberg.spark.SparkCatalog, \ spark.sql.catalog.CATALOG_NAME.catalog-impl=org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog, \ spark.sql.catalog.CATALOG_NAME.gcp_project=PROJECT_ID, \ spark.sql.catalog.CATALOG_NAME.gcp_location=LOCATION, \ spark.sql.catalog.CATALOG_NAME.warehouse=WAREHOUSE_DIRECTORY"
Substitua:
SQL_SCRIPT_PATH: o caminho para o arquivo SQL que o job em lote usa.PROJECT_ID: o ID do Google Cloud projeto em que o job em lote será executado.REGION: a região em que sua carga de trabalho é executada.SUBNET_NAME(opcional): o nome de uma sub-rede VPC naREGIONque atende aos requisitos de sub-rede da sessão.BUCKET_PATH: o local do bucket do Cloud Storage para fazer upload das dependências da carga de trabalho. OWAREHOUSE_DIRECTORYestá localizado nesse bucket. O prefixo de URIgs://do bucket não é obrigatório. É possível especificar o caminho ou o nome do bucket, por exemplo,mybucketname1.LOCATION: o local para executar o job em lote.
Para mais informações sobre como enviar jobs em lote do Apache Spark, consulte Executar uma carga de trabalho em lote do Apache Spark.
PySpark
Crie um arquivo Python com os comandos PySpark que você quer executar no catálogo de ambientes de execução do Lakehouse.
Por exemplo, o comando a seguir configura um ambiente do Apache Spark para interagir com tabelas do Apache Iceberg armazenadas no catálogo de ambientes de execução do Lakehouse. Em seguida, o comando cria um novo namespace e uma tabela do Apache Iceberg nesse namespace.
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("Lakehouse runtime catalog Iceberg") \ .config("spark.sql.catalog.CATALOG_NAME", "org.apache.iceberg.spark.SparkCatalog") \ .config("spark.sql.catalog.CATALOG_NAME.catalog-impl", "org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog") \ .config("spark.sql.catalog.CATALOG_NAME.gcp_project", "PROJECT_ID") \ .config("spark.sql.catalog.CATALOG_NAME.gcp_location", "LOCATION") \ .config("spark.sql.catalog.CATALOG_NAME.warehouse", "WAREHOUSE_DIRECTORY") \ .getOrCreate() spark.sql("USE `CATALOG_NAME`;") spark.sql("CREATE NAMESPACE IF NOT EXISTS NAMESPACE_NAME;") spark.sql("USE NAMESPACE_NAME;") spark.sql("CREATE TABLE TABLE_NAME (id int, data string) USING ICEBERG LOCATION 'WAREHOUSE_DIRECTORY';")
Substitua:
PROJECT_ID: o ID do Google Cloud projeto em que o job em lote será executado.LOCATION: o local dos recursos do BigQuery.CATALOG_NAME: o nome do catálogo que faz referência à tabela do Apache Spark.TABLE_NAME: um nome para a tabela do Apache Spark.WAREHOUSE_DIRECTORY: o URI da pasta do Cloud Storage em que o data warehouse está armazenado.NAMESPACE_NAME: o nome do namespace que faz referência à tabela do Apache Spark.
Envie o job em lote usando o seguinte
gcloud dataproc batches submit pysparkcomando:gcloud dataproc batches submit pyspark PYTHON_SCRIPT_PATH \ --version=2.2 \ --project=PROJECT_ID \ --region=REGION \ --deps-bucket=BUCKET_PATH \ --properties="spark.sql.catalog.CATALOG_NAME=org.apache.iceberg.spark.SparkCatalog,spark.sql.catalog.CATALOG_NAME.catalog-impl=org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog,spark.sql.catalog.CATALOG_NAME.gcp_project=PROJECT_ID,spark.sql.catalog.CATALOG_NAME.gcp_location=LOCATION,spark.sql.catalog.CATALOG_NAME.warehouse=WAREHOUSE_DIRECTORY"
Substitua:
PYTHON_SCRIPT_PATH: o caminho para o script Python que o job em lote usa.PROJECT_ID: o ID do Google Cloud projeto em que o job em lote será executado.REGION: a região em que sua carga de trabalho é executada.BUCKET_PATH: o local do bucket do Cloud Storage para fazer upload das dependências da carga de trabalho. O prefixo de URIgs://do bucket não é obrigatório. É possível especificar o caminho ou o nome do bucket, por exemplo,mybucketname1.
Para mais informações sobre como enviar jobs em lote do PySpark, consulte a referência do PySpark gcloud.
A seguir
- Criar e gerenciar recursos.
- Configure recursos opcionais do catálogo de ambientes de execução do Lakehouse.