Tous les produits
Search
Centre de documentation

E-MapReduce:Utiliser Spark pour écrire des données dans une table Iceberg et les lire en mode batch

Dernière mise à jour :Aug 09, 2026

Utilisez l'API DataFrame de Spark pour écrire des données dans une table Iceberg et les lire en mode batch sur un cluster Hadoop E-MapReduce (EMR). Cette rubrique utilise Spark 3.x.

Prérequis

Avant de commencer, vérifiez que vous disposez des éléments suivants :

  • Un cluster EMR Hadoop exécutant EMR V3.38.0, EMR V5.4.0 ou une version mineure ultérieure. Pour plus d'informations, consultez la section Créer un cluster.

Étape 1 : Ajouter les dépendances Maven

Créez un projet Maven et ajoutez les dépendances Spark et Iceberg au fichier POM (Project Object Model).

L'exemple suivant ajoute les dépendances pour Spark 3.1.1 et Iceberg 0.12.0 avec la portée provided. Les dépendances sont ainsi compilées localement mais résolues à partir du cluster lors de l'exécution :

<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>
Le package Iceberg inclus dans votre cluster EMR diffère de la dépendance Iceberg open source. Par exemple, l'intégration du catalogue Data Lake Formation (DLF) est native. Utilisez la dépendance Iceberg open source avec la portée provided pour compiler et empaqueter le code localement, puis laissez le package intégré du cluster prendre le relais lors de l'exécution.

Étape 2 : Configurer un catalogue

Avant d'appeler l'API DataFrame de Spark sur une table Iceberg, configurez un catalogue dans l'objet SparkConf. Le nom du catalogue et les paramètres requis varient selon la version du cluster. Tous les exemples ci-dessous utilisent DLF pour gérer les métadonnées. Pour plus de détails sur les paramètres, consultez la section Configuration des métadonnées DLF.

Sélectionnez la configuration correspondant à la version de votre cluster :

  • EMR V3.40 ou ultérieur, EMR V5.6.0 ou ultérieur — nom du catalogue : 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 — nom du catalogue : 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 — nom du catalogue : dlf_catalog

    Définissez les variables d'environnement ALIBABA_CLOUD_ACCESS_KEY_ID et ALIBABA_CLOUD_ACCESS_KEY_SECRET avant d'exécuter le code. Pour obtenir des instructions, consultez la section Configurer les variables d'environnement .
    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>")

Étape 3 : Écrire des données dans une table Iceberg

Utilisez l'API DataFrameWriterV2 (et non l'API DataFrameWriterV1) pour écrire des données depuis Spark 3.x vers une table Iceberg. L'API v2 est privilégiée car chaque opération correspond directement à une instruction SQL et le comportement d'écrasement est explicite :

Opération DataFrame Équivalent SQL Description
df.writeTo(t).create() CREATE TABLE AS SELECT Crée une nouvelle table à partir du DataFrame
df.writeTo(t).replace() REPLACE TABLE AS SELECT Remplace une table existante
df.writeTo(t).createOrReplace() CREATE OR REPLACE TABLE AS SELECT Crée la table si elle n'existe pas, ou la remplace
df.writeTo(t).append() INSERT INTO Ajoute des lignes à une table existante
df.writeTo(t).overwritePartitions() INSERT OVERWRITE dynamique Écrase les partitions correspondantes

Dans les exemples suivants, remplacez <yourCatalogName> par le nom réel de votre catalogue.

Créer une table :

val df: DataFrame = ...
df.writeTo("<yourCatalogName>.iceberg_db.sample").create()

Utilisez tableProperty() pour définir les propriétés de la table et partitionedBy() pour définir les champs de partition lors de la création de la table.

Ajouter des données :

val df: DataFrame = ...
df.writeTo("<yourCatalogName>.iceberg_db.sample").append()

Écraser les partitions :

val df: DataFrame = ...
df.writeTo("<yourCatalogName>.iceberg_db.sample").overwritePartitions()

Étape 4 : Lire des données depuis une table Iceberg

Sélectionnez la méthode de lecture en fonction de votre version de Spark :

  • Spark 3.x (recommandé)

    val df = spark.table("<yourCatalogName>.iceberg_db.sample")
  • Spark 2.4

    val df = spark.read.format("iceberg").load("<yourCatalogName>.iceberg_db.sample")

Exemple : lecture/écriture batch de bout en bout

Cet exemple utilise un cluster EMR V5.3.0 avec le catalogue dlf_catalog. Le nom du catalogue et les paramètres varient selon la version du cluster ; pour plus de détails, consultez la section Configuration des métadonnées DLF.

Étape 1 : Créer une base de données Iceberg

Utilisez Spark SQL pour créer une base de données nommée iceberg_db. Pour obtenir des instructions, consultez la section Utiliser Iceberg.

Étape 2 : Écrire l'application Spark

Le code Scala suivant configure le catalogue dlf_catalog, crée une table Iceberg, ajoute des données et lit le résultat :

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()
}

Remplacez les espaces réservés suivants par des valeurs réelles :

Espace réservé Description
<yourOSSEndpoint> Endpoint OSS, par exemple oss-cn-hangzhou.aliyuncs.com
<yourOSSWarehousePath> Chemin OSS utilisé comme racine de l'entrepôt Iceberg
<yourCatalogId> ID du catalogue DLF
<yourDLFEndpoint> Endpoint DLF
<yourDLFRegionId> ID de la région où DLF est déployé, par exemple cn-hangzhou

Étape 3 : Construire et déployer l'application

  1. Ajoutez le plugin Maven Scala au fichier pom.xml pour compiler les fichiers source 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. Déboguez le code localement, puis exécutez la commande suivante pour l'empaqueter :

    mvn clean install
  3. Connectez-vous à votre cluster EMR via SSH. Pour obtenir des instructions, consultez la section Se connecter à un cluster.

  4. Téléchargez le package JAR sur le cluster EMR. Cet exemple place le package JAR dans le répertoire racine du cluster.

Étape 4 : Soumettre le job Spark

Exécutez la commande spark-submit suivante pour soumettre le 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
Cet exemple utilise un package JAR nommé iceberg-demos.jar avec la classe principale com.aliyun.iceberg.IcebergTest . Mettez à jour --class et le nom du fichier JAR pour qu'ils correspondent à votre projet.

La sortie suivante s'affiche une fois le job terminé :

+---+----+
| id|data|
+---+----+
|  4|   d|
|  1|   a|
|  5|   e|
|  6|   f|
|  2|   b|
|  3|   c|
+---+----+

Étapes suivantes