Un nœud Flink SQL Batch vous permet de définir et d'exécuter des tâches de traitement de données à l'aide d'instructions SQL standard. Utilisez-le pour analyser et transformer de grands ensembles de données, par exemple pour le nettoyage ou l'agrégation des données. Ce nœud prend en charge la configuration visuelle et offre une solution efficace et flexible pour le traitement par lots à grande échelle. Cette rubrique explique comment utiliser un nœud Flink SQL Batch pour traiter les données par lots.
Prérequis
Créez un espace de travail et liez-y une ressource de calcul Realtime Compute for Apache Flink dans Administration. Pour plus d'informations, consultez Lier des ressources de calcul.
Créez un nœud Flink SQL Batch. Pour plus d'informations, consultez Créer un nœud pour un workflow de planification.
-
Accordez les autorisations API suivantes à l'utilisateur RAM ou au rôle RAM utilisé par DataWorks pour appeler les API Realtime Compute for Apache Flink. Ces autorisations sont requises pour soumettre et déployer les tâches du nœud sur un cluster Flink. Pour plus d'informations, consultez Accorder des autorisations.
{ "Version": "1", "Statement": [ { "Effect": "Allow", "Action": ["stream:CreateDeployment", "stream:UpdateDeployment", "stream:GetDeployment", "stream:DeleteDeployment"], "Resource": ["*"] } ] }
Limites
Seul un groupe de ressources serverless est pris en charge. L'ancien groupe de ressources exclusif pour la planification n'est pas pris en charge.
Étape 1 : Développer le nœud Flink SQL Batch
Sur la page de modification du nœud Flink SQL Batch, développez la tâche du nœud.
Développer le code SQL
Développez votre code de tâche dans la zone d'édition SQL. Dans votre code, vous pouvez définir des variables en utilisant le format ${nom_variable}. Ensuite, sur le côté droit de la page de modification du nœud, attribuez une valeur à la variable dans la section Scheduling Parameters du volet Scheduling Settings. Cela vous permet de transmettre dynamiquement des paramètres à votre code dans les scénarios de planification. Pour plus d'informations sur l'utilisation des paramètres de planification, consultez Sources des paramètres de planification et leurs expressions. Voici un exemple.
-- Créer une table source nommée datagen_source. CREATE TEMPORARY TABLE datagen_source_${var}( name VARCHAR ) WITH ( 'connector' = 'datagen', 'number-of-rows' = '1000' ); -- Créer une table de résultat nommée blackhole_sink. CREATE TEMPORARY TABLE blackhole_sink_${var}( name VARCHAR ) WITH ( 'connector' = 'blackhole' ); -- Insérer les données de la table source dans la table de résultat. INSERT INTO blackhole_sink_${var} SELECT name FROM datagen_source_${var};
Dans cet exemple, le paramètre bizdate a pour valeur $[yyyymmdd], ce qui permet la synchronisation par lot des nouvelles données quotidiennes.
Étape 2 : Configurer le nœud Flink SQL Batch
Configurez les paramètres de tâche du nœud Flink SQL Batch en fonction de vos besoins métier.
Configurer les ressources Flink
Configurez les paramètres suivants sur le côté droit de la page de modification, dans la section Flink resource information sous Scheduling Settings. Pour plus d'informations, consultez Configurer les paramètres de planification.
|**Paramètre**
|
**Description**
| | --- | --- | |
**Flink cluster**
|
Nom de la ressource de calcul Flink entièrement gérée liée dans **Administration**.
| |
**Flink engine version**
|
Sélectionnez une version du moteur en fonction de vos besoins.
| |
**Resource Group for Scheduling**
|
Sélectionnez un [groupe de ressources serverless](t2552556.xdita#) disposant d'une connectivité réseau avec Flink.
| |
**Job Manager CPU**
|
Selon les bonnes pratiques Flink, le JobManager nécessite au moins 0,5 cœur de processeur et 2 Go de mémoire pour un fonctionnement stable. Nous recommandons 1 cœur de processeur et 4 Go de mémoire, avec un maximum de 16 cœurs de processeur. Ajustez la configuration en fonction de l'échelle de votre cluster et de la complexité des tâches.
| |
**Job Manager Memory**
|
La configuration de la mémoire du JobManager affecte sa capacité à gérer les tâches de planification et de gestion. La plage recommandée est de 2 Go à 64 Go pour un fonctionnement stable et efficace. Ajustez la valeur en fonction de l'échelle de votre cluster et des exigences des tâches.
| |
**Task Manager CPU**
|
La configuration du processeur du TaskManager affecte sa capacité de traitement des tâches. Selon les bonnes pratiques Flink, nous recommandons au moins 0,5 cœur de processeur et 2 Go de mémoire, avec une configuration recommandée de 1 cœur de processeur et 4 Go de mémoire, jusqu'à un maximum de 16 cœurs de processeur. Ajustez la configuration en fonction de vos besoins.
| |
**Task Manager Memory**
|
La configuration de la mémoire du TaskManager détermine le volume de données et les performances du traitement des tâches. Pour garantir une exécution stable et efficace des tâches, la taille de la mémoire doit être d'au moins 2 Go et peut aller jusqu'à 64 Go.
| |
**Concurrency**
|
Ce paramètre détermine le nombre d'exécutions de tâches parallèles dans une tâche Flink. Une concurrence plus élevée peut améliorer la vitesse de traitement et l'utilisation des ressources. Définissez cette valeur en fonction des ressources de votre cluster et des caractéristiques de la tâche.
| |
**Maximum number of slots**
|
Un slot représente une unité de ressource de taille fixe sur un Task Manager qui peut être allouée aux tâches. Chaque slot peut exécuter une instance de tâche ou d'opérateur. Vous pouvez ajuster le nombre maximal de slots en fonction des ressources disponibles.
| |
**Number of slots per TaskManager**
|
Le nombre de slots par TaskManager détermine le nombre de tâches qu'il peut exécuter en parallèle. Vous pouvez ajuster la configuration des slots pour optimiser l'utilisation des ressources et la capacité de traitement parallèle.
|
(Facultatif) Configurer les paramètres de planification
Sur le côté droit de la page de modification, dans la section Flink cluster sous Flink engine version, cliquez sur Resource Group for Scheduling, puis modifiez le Concurrency et la Maximum number of slots pour les utiliser dynamiquement dans votre code.
(Facultatif) Configurer les paramètres d'exécution Flink
Configurez les paramètres d'exécution sur le côté droit de la page de modification, dans la section Number of slots per TaskManager sous Scheduling Parameters. Pour plus d'informations, consultez Configurer les paramètres de planification.
Lors de la configuration des paramètres d'exécution Flink, la syntaxe est compatible avec VVP (Ververica Platform). Vous pouvez écrire les configurations directement au format YAML sans ajouter de points-virgules ni d'autres caractères spéciaux pour les sauts de ligne.
Pour exécuter la tâche du nœud selon une planification périodique, configurez les informations de planification (Scheduling Settings, Add parameters, Parameter name et Parameter Value) en fonction de vos besoins métier. Pour plus d'informations, consultez Configurer les paramètres de planification.
Une fois la configuration de la tâche terminée, cliquez sur Flink running parameters.
Étape 3 : (Facultatif) Déboguer le nœud Flink SQL Batch
Avant de déployer le nœud dans l'environnement de production, utilisez la fonctionnalité de débogage pour effectuer un essai du code du nœud avec des données fictives téléchargées. Cela vous permet de vérifier la logique SQL ainsi que les données en amont et en aval sans déployer la tâche dans Operation Center.
La fonctionnalité de débogage est disponible via une liste d'autorisation. Pour l'utiliser, soumettez un ticket pour demander l'accès.
Configurer les informations sur les ressources Flink
Sur le côté droit de la page de modification du nœud, dans la section Scheduling Settings du volet Scheduling Policy, configurez les paramètres comme décrit dans le tableau suivant.
|**Paramètre**
|
**Description**
| | --- | --- | |
**Flink Debug Cluster**
|
Cluster de session Flink utilisé pour exécuter la tâche de débogage. Ce paramètre est obligatoire. La liste déroulante affiche les clusters de session existants sous la ressource de calcul actuelle ainsi que leur état d'exécution. Seuls les clusters dont l'état est **Running** peuvent être sélectionnés.
Si aucun cluster n'est disponible dans la liste, cliquez sur **Create Cluster** pour accéder à la console Realtime Compute for Apache Flink et créer un cluster de session.
| |
**Flink Engine Version**
|
Version du moteur Flink du cluster de session sélectionné. Cette valeur est affichée automatiquement par le système en fonction du cluster et ne nécessite aucune saisie manuelle.
| |
**Timeout**
|
Durée maximale d'une seule tâche de débogage, en minutes. La valeur par défaut est de 30 minutes. La tâche de débogage s'arrête automatiquement après l'écoulement de la durée spécifiée.
|
Après avoir changé la ressource de calcul pour le nœud actuel, le Scheduling time sélectionné et les données de débogage téléchargées sont effacés. Vous devez resélectionner un cluster et recharger les données.
Préparer les données de débogage
Dans la section Scheduling Dependency du volet Node output parameters, préparez des données fictives pour les tables source référencées dans votre code.
Cliquez sur Save. Le système analyse les tables source référencées dans le code SQL actuel et génère les enregistrements de noms de table correspondants dans la liste ci-dessous. Les données précédemment téléchargées ne sont pas effacées.
Dans la colonne Flink resource information de l'enregistrement du nom de table, cliquez sur Run Configuration pour télécharger un modèle CSV correspondant à la structure de la table source.
Remplissez les données de débogage dans le modèle téléchargé en respectant l'ordre des colonnes et enregistrez le fichier au format CSV.
Dans la colonne Flink Debug Cluster de l'enregistrement du nom de table, cliquez sur Running et sélectionnez le fichier CSV complété pour le charger. Une fois le chargement réussi, la colonne Create Cluster affiche Flink Engine Version.
(Facultatif) Une fois le chargement réussi, vous pouvez cliquer sur Timeout dans le panneau inférieur pour afficher les données. Pour modifier les données, rechargez un fichier CSV pour écraser les données existantes.
Si vous ne souhaitez pas que les données fictives d'une table source spécifique soient utilisées dans la session de débogage actuelle, cliquez sur Flink Debug Cluster pour changer l'état en Debug Data. Pour réactiver les données, cliquez sur Run Configuration. Seules les données dont l'état est Generate Template sont utilisées dans la session de débogage actuelle.
Avant de charger les données de débogage, vous devez d'abord sélectionner un Actions. Sinon, vous serez invité à sélectionner d'abord une ressource de calcul.
Les données de débogage prennent uniquement en charge le format CSV et la taille du fichier ne peut pas dépasser 1 Mo. La première ligne du fichier CSV doit contenir les noms des colonnes. Nous vous recommandons d'utiliser l'encodage UTF-8.
Exécuter la tâche de débogage
Une fois les données de débogage prêtes, cliquez sur le bouton Download Template dans la barre d'outils de l'éditeur (ou appuyez sur F8). Le système soumet le code, les données fictives et les informations sur les ressources Flink au cluster de session sélectionné pour exécution.
Si votre code utilise des paramètres au format ${nom_variable}, assurez-vous d'avoir attribué des valeurs aux variables dans la section Actions. Pendant le débogage, le système remplace les espaces réservés dans le code par les valeurs attribuées avant la soumission.
Afficher les résultats du débogage
Après l'exécution de la tâche de débogage, la zone de résultats en bas du nœud fournit les informations suivantes pour vous aider à identifier rapidement les problèmes :
Code : Le code SQL soumis au moteur Flink pour cette exécution (avec les substitutions de variables appliquées).
Logs : Le journal d'exécution et les informations d'erreur de la tâche de débogage.
Query results : Les données de sortie de la tâche de débogage.
Étape 4 : Déployer et gérer le nœud Flink SQL Batch
Une fois la tâche du nœud configurée, vous devez déployer le nœud. Pour plus d'informations, consultez Déployer un nœud.
Une fois la tâche déployée, cliquez sur Upload sous Deploy to Production pour afficher l'état d'exécution des tâches planifiées dans Operation Center. Pour plus d'informations, consultez Afficher les tâches planifiées.