Tous les produits
Search
Centre de documentation

Data Lake Formation:Accéder à DLF depuis un environnement Spark sur EMR on ECS

Dernière mise à jour :Aug 24, 2026

Cette rubrique explique comment accéder à un catalogue Data Lake Formation (DLF) depuis un environnement Spark sur EMR on ECS en utilisant le service REST Paimon.

Prérequis

  • Vous avez créé un cluster EMR exécutant la version 5.12.0 ou ultérieure, avec les composants Spark3 et Paimon sélectionnés. Pour demander la prise en charge d'une autre version, contactez l'équipe de développement DLF via le groupe DingTalk ID 106575000021.

  • Vous avez suivi les étapes décrites dans Premiers pas avec DLF.

  • Le cluster EMR et DLF se trouvent dans la même région, et vous avez ajouté le VPC du cluster EMR à la liste d'autorisation de DLF.

Créer un catalogue DLF

Pour plus d'informations, consultez la section Premiers pas avec DLF.

Accorder des autorisations DLF à un rôle

  1. Accordez des autorisations RAM au rôle AliyunECSInstanceForEMRRole. Vous pouvez ignorer cette étape une fois l'intégration des produits EMR et DLF terminée.

    1. Connectez-vous à la console RAM avec votre compte Alibaba Cloud ou en tant qu'administrateur RAM.

    2. Dans le volet de navigation de gauche, accédez à Identities > Roles. Recherchez ensuite le rôle AliyunECSInstanceForEMRRole.

    3. Dans la colonne Actions, cliquez sur Add Permissions.

    4. Sous l'onglet Permission Policies, recherchez et sélectionnez la stratégie AliyunDLFFullAccess. Cliquez ensuite sur OK.

  2. Accordez des autorisations DLF au rôle AliyunECSInstanceForEMRRole.

    1. Connectez-vous à la console DLF.

    2. Sur la page Catalogs, cliquez sur le nom d'un catalogue pour accéder à sa page de détails.

    3. Pour accorder des autorisations sur l'ensemble du catalogue, cliquez sur l'onglet Permissions. Sinon, pour accorder des autorisations sur une base de données ou une table spécifique, accédez d'abord à la ressource, puis cliquez sur son onglet Permissions.

    4. Dans le panneau Grant Permissions, configurez les paramètres suivants et cliquez sur OK.

      • Principal : Sélectionnez RAM User/RAM Role.

      • Select Principal : Sélectionnez AliyunECSInstanceForEMRRole dans la liste déroulante.

        Remarque

        Si AliyunECSInstanceForEMRRole ne figure pas dans la liste déroulante, cliquez sur Sync User/Role sur la page de gestion des utilisateurs.

      • Predefined Permission Type : Personnalisez les autorisations de lecture ou sélectionnez un type prédéfini, tel que Data Reader ou Data Editor.

Mettre à niveau les dépendances Paimon d'EMR

Téléchargez les deux fichiers JAR suivants (version 1.1 ou ultérieure) depuis le référentiel Maven : paimon-jindo-*.jar et paimon-spark-3.x-*.jar. Assurez-vous que les dépendances correspondent à la version Spark de votre cluster EMR.

  1. Importez les dépendances Paimon.

    1. Téléversez les deux fichiers JAR de dépendance, paimon-jindo-*.jar et paimon-spark-3.x-*.jar, vers Object Storage Service (OSS). Définissez la liste de contrôle d'accès (ACL) des fichiers sur lecture publique. Pour plus d'informations, consultez la section Téléchargement simple.

    2. Modifiez le script suivant et téléversez-le vers OSS.

      #!/bin/bash
      echo 'clean up paimon-dlf-2.5 exists file'
      rm -rf /opt/apps/PAIMON/paimon-dlf-2.5
      rm -rf /opt/apps/PAIMON/paimon-dlf-2.5.tar.gz.*
      cd /opt/apps/PAIMON/paimon-current/lib/spark3
      mkdir -p /opt/apps/PAIMON/paimon-dlf-2.5/lib/spark3
      cd /opt/apps/PAIMON/paimon-dlf-2.5/lib/spark3
      wget ${paimon-jindo-1.1.0.jar}
      wget ${paimon-spark-3.x-1.1.0.jar}
      echo 'link paimon-current to paimon-dlf-2.5'
      rm -f /opt/apps/PAIMON/paimon-current
      ln -sf /opt/apps/PAIMON/paimon-dlf-2.5 /opt/apps/PAIMON/paimon-current
      Important

      Remplacez les espaces réservés ${paimon-jindo-1.1.0.jar} et ${paimon-spark-3.x-1.1.0.jar} dans le script par les chemins de téléchargement OSS correspondants. Par défaut, les clusters EMR on ECS n'ont pas d'accès au réseau public.

      • Endpoint interne : https://{bucket}.oss-cn-hangzhou-internal.aliyuncs.com/jars/paimon-jindo-1.1.0.jar

      • Endpoint public : https://{bucket}.oss-cn-hangzhou.aliyuncs.com/jars/paimon-jindo-1.1.0.jar

  2. Exécutez le script en tant que script d'amorçage du cluster EMR. Pour plus d'informations, consultez la section Exécuter manuellement un script.

    1. Dans votre cluster EMR, accédez à l'onglet Script Actions > Run Script, puis cliquez sur Add and Run.

    2. Dans la boîte de dialogue qui s'affiche, configurez les paramètres suivants et cliquez sur OK.

      • Name : Saisissez un nom personnalisé pour le script.

      • Script Location : Sélectionnez le script de mise à niveau que vous avez téléversé vers OSS. Le chemin du script doit être au format oss://**/*.sh.

      • Target Scope : Sélectionnez Cluster.

  3. Une fois le script terminé, redémarrez le service Spark pour appliquer les modifications.

Lire et écrire des données avec Spark

Se connecter au catalogue Paimon

Exécutez la commande spark-sql suivante dans le terminal.

Important

Dans la commande, remplacez ${regionID} par votre ID de région, par exemple cn-hangzhou. Remplacez ${catalog} par le nom de votre catalogue DLF.

spark-sql --master yarn \
  --conf spark.driver.memory=5g \
  --conf spark.sql.defaultCatalog=paimon \
  --conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \
  --conf spark.sql.catalog.paimon.metastore=rest \
  --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions \
  --conf spark.sql.catalog.paimon.uri=http://${regionID}-vpc.dlf.aliyuncs.com \
  --conf spark.sql.catalog.paimon.warehouse=${catalog} \
  --conf spark.sql.catalog.paimon.token.provider=dlf \
  --conf spark.sql.catalog.paimon.dlf.token-loader=ecs

Créer des tables

Exécutez les instructions SQL suivantes pour créer des tables.

CREATE TABLE user_samples
(
    user_id INT,             
    age INT,           
    gender_code STRING,    
    clk BOOLEAN
);
CREATE TABLE user_samples_di (
    user_id INT,             
    age INT,           
    gender_code STRING,    
    clk BOOLEAN
)
USING CSV
OPTIONS(
'path'='oss://${bucket}/user/user_samples_di'
);
Remarque
  • Si vous ne spécifiez pas de base de données, les tables de données sont créées par défaut dans la base de données default du catalogue, mais vous pouvez également créer et spécifier d'autres bases de données.

  • Le répertoire /user/user_samples doit être créé préalablement dans OSS. Lorsque vous spécifiez le path, une table externe est créée. DLF stocke et gère les métadonnées de la table, tandis que les fichiers de données restent au chemin OSS spécifié. Si vous supprimez la table, seules ses métadonnées sont effacées. Les fichiers de données dans OSS ne sont pas affectés.

Insérer des données

Exécutez les instructions SQL suivantes pour insérer des données.

INSERT INTO user_samples VALUES
(1, 25, 'M', true),
(2, 18, 'F', false);
INSERT INTO user_samples_di VALUES
(1, 25, 'M', true),
(2, 18, 'F', true),
(3, 35, 'M', true);

Interroger des données

Exécutez les instructions SQL suivantes pour interroger des données.

SELECT * FROM user_samples;
SELECT * FROM user_samples_di;

La commande renvoie les résultats suivants.

spark-sql (default)> SELECT * FROM user_samples;
1	25	M	true
2	18	F	false
spark-sql (default)> SELECT * FROM user_samples_di;
1	25	M	true
2	18	F	true
3	35	M	true

Fusionner des données

Utilisez l'instruction MERGE INTO pour fusionner la table user_samples_di dans la table user_samples :

MERGE INTO user_samples
USING user_samples_di
ON user_samples.user_id = user_samples_di.user_id
WHEN MATCHED THEN
UPDATE SET
  age = user_samples_di.age,
  gender_code = user_samples_di.gender_code,
  clk = user_samples_di.clk
WHEN NOT MATCHED THEN
  INSERT (user_id, age, gender_code, clk)
  VALUES (user_samples_di.user_id, user_samples_di.age, user_samples_di.gender_code, user_samples_di.clk);

Cette opération met à jour les lignes de user_samples avec les données de user_samples_di pour les enregistrements ayant un user_id correspondant. Elle insère également de nouvelles lignes provenant de user_samples_di pour les valeurs de user_id absentes de la table user_samples.

spark-sql (default)> SELECT * FROM user_samples;
1	25	M	true
2	18	F	true
3	35	M	true