Tous les produits
Search
Centre de documentation

E-MapReduce:Use Spark to write data to an Iceberg table in streaming mode

Dernière mise à jour :Aug 09, 2026

Utilisez Apache Spark Structured Streaming pour écrire en continu les données d'Apache Kafka vers une table Apache Iceberg sur EMR. Toutes les écritures passent par l'API DataStreamWriter.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Un cluster DataLake ou un cluster personnalisé. Pour plus d'informations, consultez la section Créer un cluster.

  • Un cluster Dataflow avec le service Kafka activé. Pour plus d'informations, consultez la section Créer un cluster.

  • Les deux clusters doivent être déployés dans le même vSwitch du même VPC.

Fonctionnement

Soumettez une tâche d'écriture en flux continu en appelant DataStreamWriter depuis Spark Structured Streaming :

val tableIdentifier: String = ...
data.writeStream
    .format("iceberg")
    .outputMode("append")
    .trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
    .option("path", tableIdentifier)
    .option("checkpointLocation", checkpointPath)
    .start()

Le paramètre tableIdentifier correspond au nom ou au chemin de la table de métadonnées cible. Deux modes de sortie sont pris en charge :

Mode Comportement Équivalent SQL
append Écrit chaque micro-lot sous forme de nouvelles lignes INSERT INTO
complete Remplace l'intégralité de la table par le dernier micro-lot INSERT OVERWRITE

Diffuser des données de Kafka vers une table Iceberg

Cet exemple lit les données d'un cluster Dataflow (Kafka) et les écrit dans une table Iceberg. Une fois le code empaqueté, soumettez-le en tant que tâche Spark à l'aide de spark-submit.

Étape 1 : Configurer un topic Kafka avec des données de test

  1. Connectez-vous au cluster Dataflow en mode SSH. Pour plus d'informations, consultez la section Se connecter à un cluster.

  2. Créez un topic nommé iceberg_test :

    kafka-topics.sh --bootstrap-server core-1-1:9092,core-1-2:9092,core-1-3:9092 \
      --topic iceberg_test \
      --partitions 3 \
      --replication-factor 2 \
      --create
  3. Produisez des données de test vers le topic :

    kafka-console-producer.sh --broker-list core-1-1:9092,core-1-2:9092,core-1-3:9092 \
      --topic iceberg_test

Étape 2 : Créer une base de données et une table dans Iceberg

Utilisez Spark SQL pour créer une base de données nommée iceberg_db et une table nommée iceberg_table. Pour plus d'informations, consultez la section Utiliser Iceberg.

Étape 3 : Configurer le projet Maven

Créez un projet Maven, ajoutez les dépendances Spark et incluez le plugin Scala Maven pour la compilation. Ajoutez les éléments suivants à votre fichier pom.xml :

<dependencies>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_2.12</artifactId>
        <version>3.1.2</version>
    </dependency>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql_2.12</artifactId>
        <version>3.1.2</version>
    </dependency>
</dependencies>
<build>
    <plugins>
        <!-- the Maven Scala plugin will compile Scala source files -->
        <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>

Étape 4 : Écrire l'application Spark

Le code Scala suivant lit les données depuis Kafka et les écrit dans une table Iceberg en mode streaming.

Important

Les paramètres du catalogue et les noms de catalogues par défaut varient selon la version du cluster. Cet exemple utilise un cluster EMR V5.3.0 avec DLF (Data Lake Formation) gérant les métadonnées sous un catalogue nommé dlf_catalog. Pour obtenir des détails sur la configuration DLF, consultez la section Configuration des métadonnées DLF.

def main(args: Array[String]): Unit = {

  // Configure the Iceberg catalog backed by DLF.
  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", "<yourAccessKeyId>")
  sparkConf.set("spark.sql.catalog.dlf_catalog.access.key.secret", "<yourAccessKeySecret>")
  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("StructuredSinkIceberg")
    .getOrCreate()

  val checkpointPath = "oss://mybucket/tmp/iceberg_table_checkpoint"
  val bootstrapServers = "192.168.XX.XX:9092"
  val topic = "iceberg_test"

  // Read from the Kafka topic on the Dataflow cluster.
  val df = spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", bootstrapServers)
    .option("subscribe", topic)
    .load()

  val resDF = df.selectExpr(
      "CAST(unbase64(CAST(key AS STRING)) AS STRING) AS strKey", // Decode Base64-encoded key to string.
      "CAST(value AS STRING) AS data")
    .select(
      col("strKey").cast(LongType).alias("id"), // Cast string key to LONG type.
      col("data")
    )

  // Write to the Iceberg table every minute.
  val query = resDF.writeStream
    .format("iceberg")
    .outputMode("append")
    .trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
    .option("path", "dlf_catalog.iceberg_db.iceberg_table")
    .option("checkpointLocation", checkpointPath)
    .start()

  query.awaitTermination()
}

Modifiez les valeurs des paramètres décrits dans le tableau suivant en fonction de vos besoins métier.

Variable Description
checkpointPath Chemin OSS pour les données de point de contrôle de Spark Structured Streaming
bootstrapServers Adresse IP privée d'un courtier Kafka dans le cluster Dataflow
topic Nom du topic Kafka auquel s'abonner

Étape 5 : Empaqueter et déployer sur EMR

  1. Après avoir testé localement, empaquetez le projet :

    mvn clean install
  2. Connectez-vous à votre cluster EMR en mode SSH. Pour plus d'informations, consultez la section Se connecter à un cluster.

  3. Téléchargez le package JAR dans le répertoire racine du cluster EMR.

Étape 6 : Soumettre la tâche Spark

Exécutez spark-submit pour démarrer la tâche de streaming :

spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --driver-memory 1g \
  --executor-cores 2 \
  --executor-memory 3g \
  --num-executors 1 \
  --packages org.apache.spark:spark-sql-kafka-0-10_2.12:<version> \
  --class com.aliyun.iceberg.StructuredSinkIceberg \
  iceberg-demos.jar
Remarque
  • Remplacez <version> par la version réelle. Le composant spark-sql-kafka doit être compatible avec vos versions de Spark et de Kafka.

  • Cet exemple utilise un package JAR nommé iceberg-demos.jar. Ajustez les paramètres --class et le nom du fichier JAR afin qu'ils correspondent à votre projet.

Pour vérifier que les données sont bien écrites, utilisez Spark SQL pour interroger les modifications. Pour plus d'informations, consultez la section Utilisation de base.

Étapes suivantes