Todos os produtos
Search
Central de documentação

E-MapReduce:Use Spark to write data to an Iceberg table and read data from the table in batch mode

Última atualização: Jun 27, 2026

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: iceberg

    sparkConf.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: dlf

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

    Defina as variáveis de ambiente ALIBABA_CLOUD_ACCESS_KEY_ID e ALIBABA_CLOUD_ACCESS_KEY_SECRET antes 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

df.writeTo(t).create()

CREATE TABLE AS SELECT

Cria uma nova tabela a partir do DataFrame

df.writeTo(t).replace()

REPLACE TABLE AS SELECT

Substitui uma tabela existente

df.writeTo(t).createOrReplace()

CREATE OR REPLACE TABLE AS SELECT

Cria a tabela caso não exista ou substitui a existente

df.writeTo(t).append()

INSERT INTO

Anexa linhas a uma tabela existente

df.writeTo(t).overwritePartitions()

INSERT OVERWRITE dinâmico

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

<yourOSSEndpoint>

Endpoint do OSS, por exemplo, oss-cn-hangzhou.aliyuncs.com

<yourOSSWarehousePath>

Caminho do OSS usado como raiz do warehouse do Iceberg

<yourCatalogId>

ID do catálogo DLF

<yourDLFEndpoint>

Endpoint do DLF

<yourDLFRegionId>

ID da região onde o DLF está implantado, por exemplo, cn-hangzhou

Etapa 3: Compilar e implantar o aplicativo

  1. 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>
  2. Depure o código localmente e execute o comando a seguir para empacotá-lo:

    mvn clean install
  3. Faça login no cluster EMR via SSH. Para obter instruções, consulte Fazer login em um cluster.

  4. 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 chamado iceberg-demos.jar com a classe principal com.aliyun.iceberg.IcebergTest . Atualize o parâmetro --class e 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