Use a API DataFrame do Spark para gravar dados em uma tabela Iceberg e ler dados dessa tabela no modo em lote em um cluster Hadoop do E-MapReduce (EMR). Este tópico aborda o uso do Spark 3.x.
Pré-requisitos
Antes de começar, verifique se você tem:
Um cluster Hadoop do EMR executando o EMR V3.38.0, EMR V5.4.0 ou uma versão secundária posterior. Para mais informações, consulte Criar um cluster.
Etapa 1: Adicionar dependências do Maven
Crie um projeto Maven e adicione as dependências do Spark e do Iceberg ao arquivo do modelo de objeto do projeto (POM).
O exemplo a seguir adiciona dependências para o Spark 3.1.1 e Iceberg 0.12.0 com escopo provided. Isso significa que as dependências são compiladas localmente, mas resolvidas pelo cluster durante a execução:
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.12</artifactId>
<version>3.1.1</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.12</artifactId>
<version>3.1.1</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-core</artifactId>
<version>0.12.0</version>
<scope>provided</scope>
</dependency>
O pacote Iceberg incluído no cluster EMR difere da dependência Iceberg de código aberto. Por exemplo, a integração com o catálogo do Data Lake Formation (DLF) já vem incorporada. Use a dependência Iceberg de código aberto com escopo provided para compilar e empacotar o código localmente. Em seguida, permita que o pacote integrado do cluster assuma o controle durante a execução.
Etapa 2: Configurar um catálogo
Antes de chamar a API DataFrame do Spark em uma tabela Iceberg, configure um catálogo no objeto SparkConf. O nome do catálogo e os parâmetros necessários variam conforme a versão do cluster. Todos os exemplos abaixo usam o DLF para gerenciar metadados. Para detalhes sobre os parâmetros, consulte Configuração de metadados do DLF.
Selecione a configuração correspondente à versão do seu cluster:
-
EMR V3.40 ou posterior, EMR V5.6.0 ou posterior — nome do catálogo:
icebergsparkConf.set("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") sparkConf.set("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog") sparkConf.set("spark.sql.catalog.iceberg.catalog-impl", "org.apache.iceberg.aliyun.dlf.hive.DlfCatalog") -
EMR V3.39.X, EMR V5.5.X — nome do catálogo:
dlfsparkConf.set("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") sparkConf.set("spark.sql.catalog.dlf", "org.apache.iceberg.spark.SparkCatalog") sparkConf.set("spark.sql.catalog.dlf.catalog-impl", "org.apache.iceberg.aliyun.dlf.hive.DlfCatalog") sparkConf.set("spark.sql.catalog.dlf.warehouse", "<yourOSSWarehousePath>") -
EMR V3.38.X, EMR V5.3.X, EMR V5.4.X — nome do catálogo:
dlf_catalogDefina as variáveis de ambiente
ALIBABA_CLOUD_ACCESS_KEY_IDeALIBABA_CLOUD_ACCESS_KEY_SECRETantes de executar o código. Para obter instruções, consulte a seção Configurar variáveis de ambiente .sparkConf.set("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") sparkConf.set("spark.sql.catalog.dlf_catalog", "org.apache.iceberg.spark.SparkCatalog") sparkConf.set("spark.sql.catalog.dlf_catalog.catalog-impl", "org.apache.iceberg.aliyun.dlf.DlfCatalog") sparkConf.set("spark.sql.catalog.dlf_catalog.io-impl", "org.apache.iceberg.hadoop.HadoopFileIO") sparkConf.set("spark.sql.catalog.dlf_catalog.oss.endpoint", "<yourOSSEndpoint>") sparkConf.set("spark.sql.catalog.dlf_catalog.warehouse", "<yourOSSWarehousePath>") sparkConf.set("spark.sql.catalog.dlf_catalog.access.key.id", System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")) sparkConf.set("spark.sql.catalog.dlf_catalog.access.key.secret", System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")) sparkConf.set("spark.sql.catalog.dlf_catalog.dlf.catalog-id", "<yourCatalogId>") sparkConf.set("spark.sql.catalog.dlf_catalog.dlf.endpoint", "<yourDLFEndpoint>") sparkConf.set("spark.sql.catalog.dlf_catalog.dlf.region-id", "<yourDLFRegionId>")
Etapa 3: Gravar dados em uma tabela Iceberg
Use a API DataFrameWriterV2 (e não a API DataFrameWriterV1) para gravar dados do Spark 3.x em uma tabela Iceberg. A API v2 é preferível porque cada operação mapeia diretamente para uma instrução SQL e o comportamento de substituição é explícito:
|
Operação do DataFrame |
SQL equivalente |
Descrição |
|
|
|
Cria uma nova tabela a partir do DataFrame |
|
|
|
Substitui uma tabela existente |
|
|
|
Cria a tabela caso não exista ou substitui a existente |
|
|
|
Anexa linhas a uma tabela existente |
|
|
|
Sobrescreve as partições correspondentes |
Nos exemplos a seguir, substitua <yourCatalogName> pelo nome real do seu catálogo.
Criar uma tabela:
val df: DataFrame = ...
df.writeTo("<yourCatalogName>.iceberg_db.sample").create()
Use tableProperty() para definir propriedades da tabela e partitionedBy() para especificar campos de partição durante a criação da tabela.
Anexar dados:
val df: DataFrame = ...
df.writeTo("<yourCatalogName>.iceberg_db.sample").append()
Sobrescrever partições:
val df: DataFrame = ...
df.writeTo("<yourCatalogName>.iceberg_db.sample").overwritePartitions()
Etapa 4: Ler dados de uma tabela Iceberg
Escolha o método de leitura adequado à sua versão do Spark:
-
Spark 3.x (recomendado)
val df = spark.table("<yourCatalogName>.iceberg_db.sample") -
Spark 2.4
val df = spark.read.format("iceberg").load("<yourCatalogName>.iceberg_db.sample")
Exemplo: leitura e gravação em lote de ponta a ponta
Este exemplo usa um cluster EMR V5.3.0 com o catálogo dlf_catalog. O nome do catálogo e os parâmetros variam conforme a versão do cluster. Para mais detalhes, consulte Configuração de metadados do DLF.
Etapa 1: Criar um banco de dados Iceberg
Use o Spark SQL para criar um banco de dados chamado iceberg_db. Para obter instruções, consulte Usar o Iceberg.
Etapa 2: Escrever o aplicativo Spark
O código Scala a seguir configura o catálogo dlf_catalog, cria uma tabela Iceberg, anexa dados e lê o resultado:
def main(args: Array[String]): Unit = {
// Configure the DLF catalog
val sparkConf = new SparkConf()
sparkConf.set("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
sparkConf.set("spark.sql.catalog.dlf_catalog", "org.apache.iceberg.spark.SparkCatalog")
sparkConf.set("spark.sql.catalog.dlf_catalog.catalog-impl", "org.apache.iceberg.aliyun.dlf.DlfCatalog")
sparkConf.set("spark.sql.catalog.dlf_catalog.io-impl", "org.apache.iceberg.hadoop.HadoopFileIO")
sparkConf.set("spark.sql.catalog.dlf_catalog.oss.endpoint", "<yourOSSEndpoint>")
sparkConf.set("spark.sql.catalog.dlf_catalog.warehouse", "<yourOSSWarehousePath>")
sparkConf.set("spark.sql.catalog.dlf_catalog.access.key.id", System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
sparkConf.set("spark.sql.catalog.dlf_catalog.access.key.secret", System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"))
sparkConf.set("spark.sql.catalog.dlf_catalog.dlf.catalog-id", "<yourCatalogId>")
sparkConf.set("spark.sql.catalog.dlf_catalog.dlf.endpoint", "<yourDLFEndpoint>")
sparkConf.set("spark.sql.catalog.dlf_catalog.dlf.region-id", "<yourDLFRegionId>")
val spark = SparkSession
.builder()
.config(sparkConf)
.appName("IcebergReadWriteTest")
.getOrCreate()
// Create the Iceberg table and write the first batch of rows
val firstDF = spark.createDataFrame(Seq(
(1, "a"), (2, "b"), (3, "c")
)).toDF("id", "data")
firstDF.writeTo("dlf_catalog.iceberg_db.sample").createOrReplace()
// Append a second batch of rows
val secondDF = spark.createDataFrame(Seq(
(4, "d"), (5, "e"), (6, "f")
)).toDF("id", "data")
secondDF.writeTo("dlf_catalog.iceberg_db.sample").append()
// Read all rows from the Iceberg table
val icebergTable = spark.table("dlf_catalog.iceberg_db.sample")
icebergTable.show()
}
Substitua os seguintes espaços reservados pelos valores reais:
|
Espaço reservado |
Descrição |
|
|
Endpoint do OSS, por exemplo, |
|
|
Caminho do OSS usado como raiz do warehouse do Iceberg |
|
|
ID do catálogo DLF |
|
|
Endpoint do DLF |
|
|
ID da região onde o DLF está implantado, por exemplo, |
Etapa 3: Compilar e implantar o aplicativo
-
Adicione o plugin Scala Maven ao arquivo pom.xml para compilar os arquivos de origem Scala:
<build> <plugins> <plugin> <groupId>net.alchim31.maven</groupId> <artifactId>scala-maven-plugin</artifactId> <version>3.2.2</version> <executions> <execution> <goals> <goal>compile</goal> <goal>testCompile</goal> </goals> </execution> </executions> </plugin> </plugins> </build> -
Depure o código localmente e execute o comando a seguir para empacotá-lo:
mvn clean install Faça login no cluster EMR via SSH. Para obter instruções, consulte Fazer login em um cluster.
Carregue o pacote JAR no cluster EMR. Este exemplo coloca o pacote JAR no diretório raiz do cluster.
Etapa 4: Enviar o job do Spark
Execute o seguinte comando spark-submit para enviar o job:
spark-submit \
--master yarn \
--deploy-mode cluster \
--driver-memory 1g \
--executor-cores 1 \
--executor-memory 1g \
--num-executors 1 \
--class com.aliyun.iceberg.IcebergTest \
iceberg-demos.jar
Este exemplo usa um pacote JAR chamadoiceberg-demos.jarcom a classe principalcom.aliyun.iceberg.IcebergTest. Atualize o parâmetro--classe o nome do arquivo JAR para corresponder ao seu projeto.
A seguinte saída é retornada após a conclusão do job:
+---+----+
| id|data|
+---+----+
| 4| d|
| 1| a|
| 5| e|
| 6| f|
| 2| b|
| 3| c|
+---+----+
Próximos passos
Usar o Iceberg — gerencie tabelas Iceberg com o Spark SQL
Configuração de metadados do DLF — referência de configuração de catálogo para todas as versões do EMR