Delta Lake est pré-intégré à E-MapReduce (EMR). Il offre à vos charges de travail Apache Spark des transactions ACID, une gestion évolutive des métadonnées ainsi qu’un traitement unifié des données par lots et en flux continu sur votre stockage de lac de données existant.
Informations de configuration
Les extensions de session Spark activent Delta Lake dans EMR. La configuration requise varie selon la version de Spark.
Spark 2.X
spark.sql.extensions io.delta.sql.DeltaSparkSessionExtension
Spark 3.X
spark.sql.extensions io.delta.sql.DeltaSparkSessionExtension
spark.sql.catalog.spark_catalog org.apache.spark.sql.delta.catalog.DeltaCatalog
Ces paramètres s’appliquent par défaut lorsque vous ajoutez le service Delta Lake à votre cluster EMR. Pour les appliquer manuellement, transmettez-les via l’option --conf au démarrage de streaming-sql.
Opérations courantes
Tous les exemples de cette rubrique utilisent une table Delta nommée delta_table.
Créer une table
CREATE TABLE delta_table (id INT) USING delta;
Insérer des données
INSERT INTO delta_table VALUES 0,1,2,3,4;
Vérifiez l’insertion :
SELECT * FROM delta_table;
Remplacer des données
INSERT OVERWRITE TABLE delta_table VALUES 5,6,7,8,9;
Vérifiez le résultat :
SELECT * FROM delta_table;
Mettre à jour des données
Ajoutez 100 à tous les ID pairs :
UPDATE delta_table SET id = id + 100 WHERE mod(id, 2) = 0;
Vérifiez la mise à jour :
SELECT * FROM delta_table;
Supprimer des données
Supprimez tous les enregistrements dont l’ID est pair :
DELETE FROM delta_table WHERE mod(id, 2) = 0;
Vérifiez la suppression :
SELECT * FROM delta_table;
Fusionner des données
L’instruction MERGE INTO fusionne les données d’une table source dans une table Delta (upsert). Elle permet d’appliquer des mises à jour conditionnelles et des insertions en une seule opération atomique.
-
Créez une table source pour l’opération de fusion.
CREATE TABLE newData(id INT) USING delta; -
Insérez des données dans la table source.
INSERT INTO newData VALUES 0,1,2,3,4,5,6,7,8,9; -
Fusionnez les données de
newDatadansdelta_table. Si l’ID d’un enregistrement denewDatacorrespond à un ID dedelta_table, ajoutez 100 à l’ID correspondant. Sinon, insérez l’enregistrement tel quel.MERGE INTO delta_table AS target USING newData AS source ON target.id = source.id WHEN MATCHED THEN UPDATE SET target.id = source.id + 100 WHEN NOT MATCHED THEN INSERT *; -
Vérifiez le résultat de la fusion.
SELECT * FROM delta_table;
Lire une table Delta en mode streaming
Cette procédure utilise l’interface CLI streaming-sql d’EMR pour configurer un pipeline de streaming à partir d’une table Delta existante.
L’étape 6 nécessite une seconde session client streaming-sql active simultanément. Ouvrez un second terminal avant de commencer l’étape 6.
Prérequis
Avant de commencer, vérifiez que vous disposez des éléments suivants :
Un cluster EMR accessible via SSH
Une table Delta existante nommée
delta_table
Configurer le pipeline de streaming
-
Démarrez le client
streaming-sql.streaming-sqlSi Delta Lake n’est pas ajouté en tant que service à votre cluster, démarrez
streaming-sqlavec le fichier JAR Delta Lake et la configuration appropriée : streaming-sql --jars /path/to/delta-core_2.11-0.6.1.jar --conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension -
Créez la table de destination.
CREATE TABLE stream_debug_table (id INT) USING DELTA; -
Créez une analyse de flux (stream scan) sur la table Delta source.
CREATE SCAN stream_delta_table on delta_table USING STREAM; -
Démarrez la tâche de streaming pour écrire les données de
delta_tableversstream_debug_table.CREATE STREAM job options ( triggerType='ProcessingTime', checkpointLocation = '/tmp/streaming_read_cp' ) INSERT INTO stream_debug_table SELECT * FROM stream_delta_table;
Vérifier le comportement du streaming
-
Dans une seconde session client
streaming-sql, interrogez l’état actuel de la table de destination.SELECT * FROM stream_debug_table; -
Dans la seconde session, insérez des données dans la table source.
INSERT INTO delta_table VALUES 801, 802; -
Interrogez à nouveau la table de destination pour confirmer la diffusion des nouvelles données en streaming.
SELECT * FROM stream_debug_table; -
Insérez un autre lot de données dans la table source.
INSERT INTO delta_table VALUES 901, 902; -
Interrogez la table de destination pour confirmer la diffusion du second lot en streaming.
SELECT * FROM stream_debug_table;