E-MapReduce permet de lire et d'écrire des données dans Paimon à l'aide de Flink SQL. Cette rubrique présente des exemples de création de catalogue, de lecture et d'écriture en flux continu (streaming), ainsi que d'exécution de requêtes OLAP.
Prérequis
Vous avez créé un cluster Dataflow ou personnalisé en sélectionnant Flink et Paimon. Pour plus d'informations, consultez la rubrique Créer un cluster.
Pour utiliser un catalogue Hive, vous devez créer un cluster personnalisé en sélectionnant Flink, Paimon et Hive. Définissez également le type de Metadata Storage Method sur Self-managed RDS ou Built-in MySQL.
Limites
Les versions E-MapReduce V3.46.0 et V5.17.0 ne prennent pas en charge les catalogues DLF et Hive.
-
Vous pouvez utiliser Flink SQL pour lire et écrire dans Paimon sur les clusters exécutant les versions E-MapReduce V3.46.0 à V3.50.X et E-MapReduce V5.12.0 à V5.16.X.
RemarquePour les versions E-MapReduce V3.51.X et ultérieures, ainsi que E-MapReduce V5.17.X et ultérieures, reportez-vous à la documentation Apache Paimon et configurez l'intégration directement dans votre cluster EMR.
Procédure
Étape 1 : Configurer les dépendances
Vous pouvez lire et écrire dans Paimon en utilisant un catalogue Filesystem, un catalogue Hive ou un catalogue DLF. Configurez les dépendances selon la méthode choisie.
Catalogue Filesystem
cp /opt/apps/PAIMON/paimon-current/lib/flink/*.jar /opt/apps/FLINK/flink-current/lib/
Catalogue Hive
cp /opt/apps/PAIMON/paimon-current/lib/flink/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/FLINK/flink-current/opt/catalogs/hive-2.3.6/*.jar /opt/apps/FLINK/flink-current/lib/
Catalogue DLF
cp /opt/apps/PAIMON/paimon-current/lib/flink/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/PAIMON/paimon-current/lib/jackson/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/METASTORE/metastore-*/hive2/*.jar /opt/apps/FLINK/flink-current/lib/
cp /opt/apps/FLINK/flink-current/opt/catalogs/hive-2.3.6/*.jar /opt/apps/FLINK/flink-current/lib/
Étape 2 : Démarrer le cluster
Cette rubrique utilise le mode session à titre d'exemple. Pour les autres modes, consultez la section Utilisation de base.
Exécutez la commande suivante pour démarrer une session YARN détachée :
yarn-session.sh --detached
Étape 3 : Créer un catalogue
Paimon stocke les données et les métadonnées dans un système de fichiers tel que HDFS ou dans un stockage objet comme OSS-HDFS. Le paramètre warehouse spécifie le chemin racine. Si le chemin du warehouse spécifié n'existe pas, Paimon le crée automatiquement. S'il existe déjà, le catalogue permet d'accéder aux tables situées à cet emplacement.
Vous pouvez également synchroniser les métadonnées vers Hive ou DLF afin de permettre à d'autres services d'accéder aux données Paimon.
Les versions E-MapReduce V3.46.0 et V5.17.0 ne prennent pas en charge les catalogues DLF et Hive.
Catalogue Filesystem
Un catalogue Filesystem stocke les métadonnées uniquement dans un système de fichiers ou un stockage objet.
-
Exécutez la commande suivante pour démarrer le client Flink SQL.
sql-client.sh -
Exécutez l'instruction Flink SQL suivante pour créer un catalogue Filesystem.
CREATE CATALOG test_catalog WITH ( 'type' = 'paimon', 'metastore' = 'filesystem', 'warehouse' = 'oss://<yourBucketName>/warehouse' );
Catalogue Hive
Un catalogue Hive synchronise les métadonnées vers Hive Metastore. Les tables créées dans un catalogue Hive peuvent être interrogées directement depuis Hive.
Pour plus de détails sur l'interrogation de Paimon depuis Hive, consultez la rubrique Intégrer Paimon avec Hive.
-
Exécutez la commande suivante pour démarrer le client Flink SQL.
sql-client.shRemarqueLa commande de démarrage est identique, quelle que soit la version de Hive.
-
Exécutez l'instruction Flink SQL suivante pour créer un catalogue Hive.
CREATE CATALOG test_catalog WITH ( 'type' = 'paimon', 'metastore' = 'hive', 'uri' = 'thrift://master-1-1:9083', -- The uri parameter specifies the address of the Hive metastore service. 'warehouse' = 'oss://<yourBucketName>/warehouse' );
Catalogue DLF
Un catalogue DLF synchronise les métadonnées vers DLF.
Lors de la création du cluster, définissez le paramètre Metadata sur DLF Unified Metadata.
-
Exécutez la commande suivante pour démarrer le client Flink SQL.
sql-client.shRemarqueLa commande de démarrage est identique, quelle que soit la version de Hive.
-
Exécutez l'instruction Flink SQL suivante pour créer un catalogue DLF.
CREATE CATALOG test_catalog WITH ( 'type' = 'paimon', 'metastore' = 'dlf', 'hive-conf-dir' = '/etc/taihao-apps/flink-conf', 'warehouse' = 'oss://<yourBucketName>/warehouse' );
Étape 4 : Lire et écrire dans Paimon via le streaming
Exécutez les instructions Flink SQL suivantes pour créer une table dans le catalogue, puis pour y lire et y écrire des données.
-- Set the execution mode to streaming.
SET 'execution.runtime-mode' = 'streaming';
-- Paimon requires you to set a checkpoint interval for streaming jobs.
SET 'execution.checkpointing.interval' = '10s';
-- Use the catalog created in the previous step.
USE CATALOG test_catalog;
-- Create and use a test database.
CREATE DATABASE test_db;
USE test_db;
-- Use datagen to generate random data.
CREATE TEMPORARY TABLE datagen_source (
uuid int,
kind int,
price int
) WITH (
'connector' = 'datagen',
'fields.kind.min' = '0',
'fields.kind.max' = '9',
'rows-per-second' = '10'
);
-- Create a Paimon table.
CREATE TABLE test_tbl (
uuid int,
kind int,
price int,
PRIMARY KEY (uuid) NOT ENFORCED
);
-- Write data to the Paimon table.
INSERT INTO test_tbl SELECT * FROM datagen_source;
-- Read data from the table.
-- The preceding streaming write job runs concurrently.
-- Ensure your Flink cluster has sufficient resources (task slots) for both jobs. Otherwise, this query will not execute.
SELECT kind, SUM(price) FROM test_tbl GROUP BY kind;
Étape 5 : Exécuter une requête OLAP sur Paimon
Exécutez les instructions Flink SQL suivantes pour effectuer une requête OLAP sur la table créée.
-- Set the execution mode to batch.
RESET 'execution.checkpointing.interval';
SET 'execution.runtime-mode' = 'batch';
-- Use tableau mode to print results directly in the terminal.
SET 'sql-client.execution.result-mode' = 'tableau';
-- Query data in the table.
SELECT kind, SUM(price) FROM test_tbl GROUP BY kind;
Étape 6 : Nettoyer les ressources
Une fois les tests terminés, arrêtez le job d'écriture en flux continu Paimon afin d'éviter toute fuite de ressources.
Après avoir arrêté le job, exécutez l'instruction Flink SQL suivante pour supprimer la table créée.
DROP TABLE test_tbl;