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 :
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 — nom du catalogue :
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 — nom du catalogue :
dlf_catalogDéfinissez les variables d'environnement
ALIBABA_CLOUD_ACCESS_KEY_IDetALIBABA_CLOUD_ACCESS_KEY_SECRETavant 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
-
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> -
Déboguez le code localement, puis exécutez la commande suivante pour l'empaqueter :
mvn clean install Connectez-vous à votre cluster EMR via SSH. Pour obtenir des instructions, consultez la section Se connecter à un cluster.
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.jaravec la classe principalecom.aliyun.iceberg.IcebergTest. Mettez à jour--classet 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
Utiliser Iceberg — gérer les tables Iceberg avec Spark SQL
Configuration des métadonnées DLF — référence de configuration du catalogue pour toutes les versions EMR