Conecte o EMR Serverless Spark a um catálogo do Data Lake Formation (DLF) usando o protocolo Iceberg REST para ler e gravar tabelas Iceberg por meio do Spark SQL.
Pré-requisitos
Antes de começar, certifique-se de ter:
Um workspace do Serverless Spark criado na mesma região da sua instância do DLF. Consulte Criar um workspace
Tipos de tarefa compatíveis
Todos os três tipos de tarefa permitem conectar ao DLF por meio do catálogo Iceberg REST:
|
Tipo de tarefa |
Referência |
|
Sessão SQL |
|
|
Spark Thrift Server |
|
|
Job em lote |
Etapa 1: Conceda permissões ao catálogo
Acesse o console do Data Lake Formation.
Na página Catalogs, clique no nome do catálogo para abrir sua página de detalhes.
Clique na guia Permissions para conceder acesso a todo o catálogo. Para conceder acesso a um banco de dados ou a uma tabela específica, acesse o recurso e clique na guia Permissions.
-
Configure os seguintes campos e clique em OK:
NotaSe AliyunECSInstanceForEMRRole não aparecer na lista suspensa, acesse a página de gerenciamento de usuários e clique em Sync.
Campo
Valor
User/Role
Selecione RAM User/RAM Role
Select Authorization Object
Selecione AliyunECSInstanceForEMRRole na lista suspensa
Preset Permission Type
Selecione as permissões de leitura manualmente ou escolha uma função predefinida, como Data Reader ou Data Editor
Se você for um usuário do Resource Access Management (RAM), conceda as permissões de recurso necessárias antes de executar operações de dados. Consulte Gerenciamento de autorização de dados.
Etapa 2: Conecte-se ao catálogo e leia ou grave dados
Escolha um dos seguintes métodos de conexão dependendo de como você gerencia o catálogo.
Opção 1: Usar um catálogo de dados (recomendado)
Se você trabalhar com um catálogo de dados gerenciado pelo DLF, não será necessário configurar a sessão do Spark. Acesse a página Data Catalog, clique em Add data catalog e selecione o catálogo diretamente no desenvolvimento em Spark SQL.
Opção 2: Usar um catálogo personalizado
Adicione a seguinte configuração na seção Spark Configuration em Custom Configuration.
A configuração abaixo usa iceberg_catalog como nome do catálogo. Esse nome registra um serviço de gerenciamento de tabelas Iceberg no Spark, suportado pelo catálogo Iceberg REST, o qual se conecta ao DLF por meio de APIs REST. Altere o nome do catálogo e os parâmetros relacionados conforme necessário.
# Enable the Iceberg Spark extension
spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
# Register the catalog
spark.sql.catalog.iceberg_catalog org.apache.iceberg.spark.SparkCatalog
# Use the Iceberg REST catalog implementation
spark.sql.catalog.iceberg_catalog.catalog-impl org.apache.iceberg.rest.RESTCatalog
# DLF Iceberg REST endpoint
spark.sql.catalog.iceberg_catalog.uri http://${regionID}-vpc.dlf.aliyuncs.com/iceberg
# Your DLF catalog name
spark.sql.catalog.iceberg_catalog.warehouse ${catalogName}
# Use the DLF FileIO implementation
spark.sql.catalog.iceberg_catalog.io-impl org.apache.iceberg.rest.DlfFileIO
# Enable SigV4 signature authentication
spark.sql.catalog.iceberg_catalog.rest.auth.type sigv4
spark.sql.catalog.iceberg_catalog.rest.auth.sigv4.delegate-auth-type none
spark.sql.catalog.iceberg_catalog.rest.signing-region ${regionID}
spark.sql.catalog.iceberg_catalog.rest.signing-name DlfNext
# Access credentials
spark.sql.catalog.iceberg_catalog.rest.access-key-id ${access_key_id}
spark.sql.catalog.iceberg_catalog.rest.secret-access-key ${access_key_secret}
Substitua os seguintes placeholders pelos seus valores reais:
|
Placeholder |
Descrição |
Exemplo |
|
|
ID da região onde sua instância do DLF está implantada. Consulte Endpoints. |
|
|
|
O nome do seu catálogo DLF |
|
|
|
AccessKey ID da sua conta Alibaba Cloud |
— |
|
|
AccessKey Secret da sua conta Alibaba Cloud |
— |
Para sessões SQL, use a versão do mecanismo esr-4.7.0, esr-3.6.0 ou posterior.
Leitura e gravação de dados
Os exemplos a seguir apresentam operações comuns do Spark SQL em uma tabela Iceberg no DLF. Todas as instruções referenciam tabelas no formato iceberg_catalog.<database>.<table>.
Se você não especificar um banco de dados, o Spark criará as tabelas no banco de dados default do catálogo.
Para um guia passo a passo completo de desenvolvimento com Spark SQL, consulte Primeiros passos com o desenvolvimento Spark SQL.
-- Create a database
CREATE DATABASE IF NOT EXISTS db;
-- Create a non-partitioned table
CREATE TABLE iceberg_catalog.db.tbl (
id BIGINT NOT NULL COMMENT 'unique id',
data STRING
)
USING iceberg;
-- Insert rows
INSERT INTO iceberg_catalog.db.tbl VALUES
(1, 'Alice'),
(2, 'Bob'),
(3, 'Charlie');
-- Query all rows
SELECT * FROM iceberg_catalog.db.tbl;
-- Query by condition
SELECT * FROM iceberg_catalog.db.tbl WHERE id = 2;
-- Update a row
UPDATE iceberg_catalog.db.tbl SET data = 'David' WHERE id = 3;
-- Confirm the update
SELECT * FROM iceberg_catalog.db.tbl WHERE id = 3;
-- Delete a row
DELETE FROM iceberg_catalog.db.tbl WHERE id = 1;
-- Confirm the deletion
SELECT * FROM iceberg_catalog.db.tbl;
-- Create a partitioned table
CREATE TABLE iceberg_catalog.db.part_tbl (
id BIGINT,
data STRING,
category STRING,
ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (category);
-- Insert rows
INSERT INTO iceberg_catalog.db.part_tbl VALUES
(100, 'Data1', 'A', to_timestamp('2025-01-01 12:00:00')),
(200, 'Data2', 'B', to_timestamp('2025-01-02 14:00:00')),
(300, 'Data3', 'A', to_timestamp('2025-01-01 15:00:00')),
(400, 'Data4', 'C', to_timestamp('2025-01-03 10:00:00'));
-- Query all rows
SELECT * FROM iceberg_catalog.db.part_tbl;
-- Filter by bucket
SELECT * FROM iceberg_catalog.db.part_tbl WHERE bucket(16, id) = 0;
-- Filter by day
SELECT * FROM iceberg_catalog.db.part_tbl WHERE days(ts) = '2025-01-01';
-- Filter by partition column
SELECT * FROM iceberg_catalog.db.part_tbl WHERE category = 'A';
-- Combined filter (bucket + day + category)
SELECT * FROM iceberg_catalog.db.part_tbl
WHERE bucket(16, id) = 0
AND days(ts) = '2025-01-01'
AND category = 'A';
-- Aggregate by category
SELECT category, COUNT(*) AS count
FROM iceberg_catalog.db.part_tbl
GROUP BY category;
-- Drop the database (all tables must be empty first)
-- DROP DATABASE iceberg_catalog.db;