Realtime Compute for Apache Flink permet de créer des tâches par lots en utilisant le dialecte Hive. Cette fonctionnalité garantit la compatibilité avec la syntaxe Hive SQL, améliore l'interopérabilité et simplifie la migration des tâches Hive existantes vers la console.
Prérequis
Si vous accédez à la console en tant qu'utilisateur RAM, rôle RAM ou autre identité, assurez-vous de disposer des autorisations requises. Pour plus d'informations, consultez la rubrique Autorisations.
Un espace de travail doit être créé. Pour plus d'informations, consultez la rubrique Activer Realtime Compute for Apache Flink.
Limites
Seules les versions Ververica Runtime (VVR) 8.0.11 et ultérieures prennent en charge le dialecte Hive.
Les tâches SQL ne prennent actuellement en charge que la syntaxe des instructions INSERT du dialecte Hive. Vous devez déclarer
USE Catalog <yourHiveCatalog>avant les instructions INSERT. Pour créer une table, effectuez l'opération sur la page Scripts.Les fonctions définies par l'utilisateur (UDF) Hive et Flink ne sont pas prises en charge.
Étape 1 : Créer un catalogue Hive
Configurez les métadonnées Hive. Pour plus d'informations, consultez la rubrique Configurer les métadonnées Hive.
-
Créez un catalogue Hive. Pour plus d'informations, consultez la rubrique Créer un catalogue Hive.
Dans ce tutoriel, le catalogue Hive est nommé
hdfshive.
Étape 2 : Préparer des tables Hive d'exemple
Dans le volet de navigation de gauche, accédez à . Cliquez sur
New pour créer un script.-
Exécutez les instructions SQL d'exemple suivantes.
ImportantLa table source Hive et la table de destination doivent être des tables permanentes créées avec l'instruction
CREATE TABLE. L'utilisation de tables temporaires créées avec l'instructionCREATE TEMPORARY TABLEn'est pas autorisée.-- Use the Hive catalog. In this example, the catalog is named hdfshive and was created in Step 1. USE CATALOG hdfshive; -- Create a source table with the default storage format. CREATE TABLE source_table ( id INT, name STRING, age INT, city STRING, salary FLOAT )WITH ('connector' = 'hive'); -- Create a sink table with the default storage format. CREATE TABLE target_table ( city STRING, avg_salary FLOAT, user_count INT )WITH ('connector' = 'hive'); -- Insert sample data into the source table. INSERT INTO source_table VALUES (1, 'Alice', 25, 'New York', 5000.0), (2, 'Bob', 30, 'San Francisco', 6000.0), (3, 'Charlie', 35, 'New York', 7000.0), (4, 'David', 40, 'San Francisco', 8000.0), (5, 'Eva', 45, 'Los Angeles', 9000.0); -- Create a table with a specific storage format, for example, Parquet. -- Load the Hive module. load MODULE hive with ('hive-version' = '2.3.6'); use CATALOG `hdfshive`; -- Required: Set the SQL dialect to 'hive' to recognize Hive DDL keywords such as 'STORED'. set 'table.sql-dialect' = 'hive'; CREATE TABLE `parquet_table`( id INT, name STRING, age INT, city STRING, salary FLOAT )STORED AS PARQUET;
Étape 3 : Créer une tâche Hive SQL
Dans le volet de navigation de gauche, accédez à .
Cliquez sur New. Dans la boîte de dialogue New Draft, sélectionnez Blank Batch Draft (BETA) et cliquez sur Next.
-
Saisissez les informations relatives à la tâche.
Paramètre
Description
Exemple
Name
Nom de la tâche.
RemarqueLe nom de la tâche doit être unique au sein de l'espace de travail actuel.
hive-sql
Location
Dossier dans lequel le fichier de code de la tâche est stocké.
Vous pouvez également cliquer sur l'icône
située à droite d'un dossier existant pour créer un sous-dossier.Drafts
Engine version
Version du moteur Flink utilisée par la tâche.
Nous vous recommandons de sélectionner une version portant le tag RECOMMENDED. Ces versions offrent une fiabilité et des performances supérieures. Pour plus d'informations sur les versions du moteur, consultez les rubriques Notes de version et Versions du moteur.
vvr-8.0.11-flink-1.17
SQL dialect
Langage SQL utilisé pour le traitement des données.
RemarqueCe paramètre n'apparaît que si vous sélectionnez une version du moteur prenant en charge le dialecte Hive.
Hive SQL
Cliquez sur Create.
Étape 4 : Rédiger et déployer la tâche Hive SQL
-
Rédigez les instructions SQL.
Cet exemple calcule le nombre d'utilisateurs âgés de plus de 30 ans et le salaire moyen pour chaque ville. Copiez le script SQL suivant dans l'éditeur SQL.
-- Use the Hive catalog. In this example, the catalog is named hdfshive and was created in Step 1. USE CATALOG hdfshive; INSERT INTO TABLE target_table SELECT city, AVG(salary) AS avg_salary, -- Calculate the average salary COUNT(id) AS user_count -- Count the number of users FROM source_table WHERE age > 30 -- Filter for users older than 30 GROUP BY city; -- Group by city Dans le coin supérieur droit, cliquez sur Deploy. Dans la boîte de dialogue, configurez les paramètres selon vos besoins (ce tutoriel utilise les paramètres par défaut) et cliquez sur OK.
(Facultatif) Étape 5 : Configurer les paramètres d'exécution
Cette étape n'est nécessaire que si vous utilisez JindoSDK pour accéder à votre cluster Hive.
Dans le volet de navigation de gauche, accédez à .
Dans la liste déroulante, sélectionnez BATCH. Recherchez la tâche cible et cliquez sur Details dans la colonne Actions.
Dans le panneau des détails du déploiement, cliquez sur Edit dans la section Runtime parameters configuration.
-
Dans le champ Other Configuration, ajoutez la configuration suivante :
fs.oss.jindo.endpoint: <YOUR_Endpoint> fs.oss.jindo.buckets: <YOUR_Buckets> fs.oss.jindo.accessKeyId: <YOUR_AccessKeyId> fs.oss.jindo.accessKeySecret: <YOUR_AccessKeySecret>Pour plus d'informations sur ces paramètres, consultez la rubrique Écrire des données dans OSS-HDFS.
Cliquez sur Save.
Étape 6 : Démarrer la tâche et afficher les résultats
Sur la page Deployments, sélectionnez Batch job dans le filtre, recherchez votre tâche cible (par exemple, hive-sql) et cliquez sur Start dans la colonne Actions.
-
Une fois que l'état de la tâche passe à FINISHED, consultez les résultats.
Sur la page , exécutez l'instruction SQL suivante pour afficher les données, qui incluent le nombre d'utilisateurs âgés de plus de 30 ans et leur salaire moyen dans chaque ville.
-- Use the Hive catalog. In this example, the catalog is named hdfshive and was created in Step 1. USE CATALOG hdfshive; select * from target_table;La requête renvoie trois lignes de target_table avec les colonnes city, avg_salary et user_count : Los Angeles (9000.0, 1), New York (7000.0, 1) et San Francisco (8000.0, 1).
Développement de tâches JAR Hive
Vous pouvez exécuter des tâches utilisant le dialecte Hive en tant que tâches JAR. Cela nécessite la version 11.2 ou ultérieure du package JAR « ververica-connector-hive-2.3.6 ». Assurez-vous également que les configurations Hive de votre tâche JAR correspondent aux paramètres de la console.
-
Paramètres de la console
L'URI JAR spécifie le package JAR téléchargé pour la tâche JAR.
Dans Additional Dependencies, téléchargez les quatre fichiers de configuration de votre cluster Hive : core-site.xml, mapred-site.xml, hdfs-site.xml et hive-site.xml. Vous devez également télécharger le package JAR ververica-connector-hive-2.3.6.
-
Configurez les paramètres d'exécution. En fonction de la configuration de votre cluster Hive, si vous devez écrire des données dans OSS-HDFS, utilisez les paramètres décrits dans la rubrique (Facultatif) Étape 5 : Configurer les paramètres d'exécution.
table.sql-dialect: HIVE classloader.parent-first-patterns.additional: org.apache.hadoop;org.antlr.runtime kubernetes.application-mode.classpath.include-user-jar: true
-
Exemple de code de tâche JAR :
-
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env); Configuration conf = new Configuration(); conf.setString("type", "hive"); conf.setString("default-database", "default"); conf.setString("hive-version", "2.3.6"); conf.setString("hive-conf-dir", "/flink/usrlib/" ); conf.setString("hadoop-conf-dir", "/flink/usrlib/"); CatalogDescriptor descriptor = CatalogDescriptor.of("hivecat", conf); tableEnv.createCatalog("hivecat", descriptor); tableEnv.loadModule("hive", new HiveModule()); tableEnv.useModules("hive"); tableEnv.useCatalog("hivecat"); tableEnv.executeSql("insert into `hivecat`.`default`.`test_write` select * from `hivecat`.`default`.`test_read`;");
-
Documentation connexe
Pour plus d'informations sur la syntaxe de l'instruction
INSERTpour le dialecte Hive, consultez la rubrique Instructions INSERT | Apache Flink.Pour plus d'informations sur le traitement des données par lots avec Flink SQL, consultez la rubrique Traitement par lots.