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
Connectez-vous au cluster Dataflow en mode SSH. Pour plus d'informations, consultez la section Se connecter à un cluster.
-
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 -
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.
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
-
Après avoir testé localement, empaquetez le projet :
mvn clean install Connectez-vous à votre cluster EMR en mode SSH. Pour plus d'informations, consultez la section Se connecter à un cluster.
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
Remplacez
<version>par la version réelle. Le composantspark-sql-kafkadoit être compatible avec vos versions de Spark et de Kafka.Cet exemple utilise un package JAR nommé
iceberg-demos.jar. Ajustez les paramètres--classet 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.