Les nœuds Flink SQL Streaming dans DataWorks Data Studio vous permettent de définir une logique de traitement en temps réel à l'aide de SQL standard. Ils prennent en charge une gestion robuste de l'état, la tolérance aux pannes, ainsi que les sémantiques d'heure d'événement et d'heure de traitement. Ils s'intègrent également à des systèmes tels que Kafka et HDFS. Cette rubrique explique comment développer, configurer et exécuter un nœud Flink SQL Streaming pour traiter des données en temps réel.
Prérequis
Une ressource de calcul pour Realtime Compute for Apache Flink doit être associée dans Administration. Pour plus d'informations, consultez Associer un moteur de calcul.
Un nœud Flink SQL Streaming a été créé. Pour plus d'informations, consultez Créer un nœud pour un workflow de planification.
-
Vous avez accordé les permissions OpenAPI requises à l'utilisateur RAM ou au rôle RAM utilisé par DataWorks pour appeler les API Realtime Compute for Apache Flink. Ces permissions permettent à DataWorks de soumettre et de déployer les tâches du nœud sur un cluster Flink.
{ "Version": "1", "Statement": [ { "Effect": "Allow", "Action": ["stream:CreateDeployment", "stream:UpdateDeployment", "stream:GetDeployment", "stream:DeleteDeployment"], "Resource": ["*"] } ] }
Limites
Ce nœud ne peut pas être utilisé dans un workflow ; il doit être développé et exécuté en tant que nœud autonome.
Seuls les groupes de ressources serverless sont pris en charge. Les anciens groupes de ressources exclusifs pour la planification ne sont pas pris en charge.
Étape 1 : Développer le nœud Flink SQL Streaming
Développez la tâche du nœud sur la page d'édition du nœud Flink SQL Streaming.
Développer le code SQL
Dans l'éditeur SQL, vous pouvez définir des variables en utilisant le format ${nom_variable}. Attribuez des valeurs à ces variables dans la section Script Parameters du panneau Real-Time configuration afin de transmettre dynamiquement des paramètres dans les scénarios de planification. Par exemple :
--Create the source table datagen_source. CREATE TEMPORARY TABLE datagen_source( name VARCHAR ) WITH ( 'connector' = 'datagen' ); --Create the result table blackhole_sink. CREATE TEMPORARY TABLE blackhole_sink( name VARCHAR ) WITH ( 'connector' = 'blackhole' ); --Insert data from the source table into the result table. INSERT INTO blackhole_sink SELECT name FROM datagen_source WHERE LENGTH(name) > ${name_length};
Dans cet exemple, la valeur du paramètre name_length est 5 . Ce paramètre filtre les données pour ne traiter que les enregistrements dont la longueur du nom dépasse 5 caractères.
Étape 2 : Configurer le nœud Flink SQL Streaming
Configurez les paramètres suivants du nœud Flink SQL Streaming en fonction de vos besoins métier.
Configurer les ressources Flink
Dans la section Flink resource information du panneau Real-Time configuration , configurez les paramètres ci-dessous selon le Resource Mode sélectionné. Pour plus d'informations, consultez Configurer les ressources Flink.
|
Paramètre |
Description |
|
Flink cluster |
La ressource de calcul Flink entièrement gérée associée dans Administration. |
|
Flink engine version |
La version du moteur à utiliser. Sélectionnez une version en fonction de vos besoins. |
|
Resource Group |
Sélectionnez un groupe de ressources serverless disposant d'une connectivité réseau avec Flink. |
|
Resource Mode prend en charge les deux modes suivants. Pour plus d'informations, consultez Configurer les ressources Flink.
Configurez les paramètres en fonction du mode de ressource que vous avez sélectionné. Une bonne compréhension de l'architecture Flink vous aidera à configurer les paramètres plus efficacement. Pour plus d'informations, consultez Flink Architecture | Apache Flink. |
|
|
Basic mode |
|
|
CPU Job Manager |
Le JobManager nécessite au moins 0,5 cœur de processeur et 2 Go de mémoire pour fonctionner de manière stable. La configuration recommandée est de 1 cœur de processeur et 4 Go de mémoire, avec un maximum de 16 cœurs de processeur. Ajustez en fonction de l'échelle du cluster et de la complexité des jobs. |
|
Mémoire Job Manager |
La mémoire du JobManager affecte la capacité de planification et de gestion. La plage recommandée est de 2 Go à 64 Go. Ajustez en fonction de l'échelle du cluster et des exigences des jobs. |
|
CPU Task Manager |
Le CPU du TaskManager influence la capacité de traitement des tâches. Au moins 0,5 cœur de processeur et 2 Go de mémoire sont recommandés, avec une configuration préférable de 1 cœur de processeur et 4 Go de mémoire. Le maximum est de 16 cœurs de processeur. Ajustez selon vos besoins. |
|
Mémoire Task Manager |
La mémoire du TaskManager détermine le volume de données et les performances de traitement. La taille de la mémoire doit être d'au moins 2 Go et peut atteindre 64 Go. |
|
Concurrency |
Le nombre d'exécutions de tâches parallèles dans un job Flink. Une concurrence plus élevée améliore la vitesse de traitement et l'utilisation des ressources. Définissez cette valeur en fonction des ressources du cluster et des caractéristiques du job. |
|
Number of slots per TaskManager |
Le nombre d'emplacements par TaskManager, qui détermine combien de tâches il peut exécuter en parallèle. Ajustez les emplacements pour optimiser l'utilisation des ressources et le traitement parallèle. |
|
Expert mode |
|
|
CPU Job Manager |
Le JobManager nécessite au moins 0,25 cœur de processeur et 1 Go de mémoire pour fonctionner de manière stable, avec un maximum de 16 cœurs de processeur. Ajustez en fonction de l'échelle du cluster et de la complexité des jobs. |
|
Mémoire Job Manager |
La mémoire du JobManager affecte la capacité de planification et de gestion. La plage recommandée est de 1 Go à 64 Go. Ajustez en fonction de l'échelle du cluster et des exigences des jobs. |
|
Number of slots per TaskManager |
Le nombre d'emplacements par TaskManager, qui détermine combien de tâches il peut exécuter en parallèle. Ajustez les emplacements pour optimiser l'utilisation des ressources et le traitement parallèle. |
|
Multiple SSG mode |
Par défaut, tous les opérateurs partagent un seul groupe de partage d'emplacements, ce qui empêche de configurer individuellement les ressources pour chaque opérateur. Activez le Multiple SSG mode pour attribuer à chaque opérateur son propre emplacement indépendant, puis configurez les ressources sur l'emplacement correspondant. |
(Facultatif) Configurer les paramètres de script
Dans la section Script Parameters du panneau Real-Time configuration situé dans le volet de navigation de droite, cliquez sur Add parameters et modifiez le Parameter name ainsi que la Parameter Value pour les utiliser dynamiquement dans votre code.
(Facultatif) Configurer les paramètres d'exécution Flink
Dans la section Flink running parameters du panneau Real-Time configuration situé dans le volet de navigation de droite, configurez les paramètres suivants. Pour plus d'informations, consultez Configurer les paramètres d'exécution Flink.
|
Paramètre |
Description |
|
System Checkpoint Interval |
L'intervalle de temps auquel Flink effectue des points de contrôle périodiques. Un intervalle plus court réduit le temps de récupération après panne, mais augmente la surcharge système. Si ce champ est laissé vide, les points de contrôle sont désactivés. |
|
Minimum time interval between two system checkpoints |
Le temps d'attente minimal entre deux points de contrôle consécutifs, empêchant des points de contrôle trop fréquents d'affecter les performances. Cela garantit un écart minimal entre deux points de contrôle lorsque le parallélisme maximal pour les points de contrôle est de 1. |
|
State Data Expiration Time |
La durée maximale pendant laquelle les données d'état peuvent être conservées sans être consultées ou mises à jour. La valeur par défaut est de 36 heures, après quoi les informations d'état expirent automatiquement et sont effacées pour optimiser l'utilisation du stockage et des ressources. Important
Cette valeur par défaut repose sur les meilleures pratiques cloud et diffère de la valeur par défaut open source (0, ce qui signifie que l'état n'expire jamais). |
|
Others |
Paramètres d'exécution Flink supplémentaires. Par exemple : Remarque
Pour plus d'informations sur la configuration des paramètres, consultez Configurer les paramètres d'exécution Flink. |
Une fois la tâche configurée, cliquez sur Save pour enregistrer la tâche du nœud.
Étape 3 : (Facultatif) Déboguer le nœud Flink SQL Streaming
Avant de déployer le nœud en production, vous pouvez le déboguer en exécutant le 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
Dans la section Flink resource information du panneau Run Configuration situé à droite de la page d'édition du nœud, configurez les paramètres comme décrit dans le tableau suivant.
|**Paramètre**
|
**Description**
| | --- | --- | |
**Flink Debug Cluster**
|
Le cluster de session Flink utilisé pour exécuter la tâche de débogage. 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**
|
La version du moteur Flink du cluster de session sélectionné. Affichée automatiquement en fonction de la sélection du cluster.
| |
**Timeout**
|
La durée d'exécution maximale d'une seule tâche de débogage, en minutes. Par défaut : 30 minutes. La tâche de débogage s'arrête automatiquement après cette durée.
|
Après avoir changé la ressource de calcul du nœud actuel, le Flink Debug Cluster sélectionné et les données de débogage téléchargées sont effacés. Vous devez resélectionner un cluster et télécharger à nouveau les données.
Préparer les données de débogage
Dans la section Running du panneau Create Cluster , préparez des données fictives pour les tables sources référencées dans votre code.
Cliquez sur Flink Engine Version . Le système analyse les tables sources référencées dans le 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 Timeout d'un enregistrement de nom de table, cliquez sur Flink Debug Cluster pour télécharger un modèle CSV correspondant au schéma de la table source.
Remplissez localement les données de débogage en respectant l'ordre des colonnes du modèle et enregistrez le fichier au format CSV.
Dans la colonne Debug Data d'un enregistrement de nom de table, cliquez sur Run Configuration et sélectionnez le fichier CSV complété pour le télécharger. Une fois le téléchargement réussi, la colonne Generate Template affiche Actions .
(Facultatif) Après un téléchargement réussi, vous pouvez cliquer sur Download Template pour afficher le contenu des données dans le panneau inférieur. Pour modifier les données, téléchargez à nouveau 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 participent à la session de débogage actuelle, cliquez sur Actions pour changer l'état en Upload . Pour les réactiver, cliquez sur Status . Seules les données dont l'état est Enabled sont utilisées lors de la session de débogage.
Avant de télécharger les données de débogage, sélectionnez d'abord un Preview . Sinon, vous serez invité à sélectionner d'abord une ressource de calcul.
Les données de débogage ne prennent en charge que le format CSV, avec une taille de fichier maximale de 1 Mo. La première ligne du fichier CSV doit contenir les noms des colonnes et l'encodage UTF-8 est recommandé.
Exécuter la tâche de débogage
Une fois les données de débogage préparées, cliquez sur le bouton Disable 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 Disabled . 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 affiche les informations suivantes :
Code : Le code SQL soumis au moteur Flink (avec substitution des variables appliquée).
Logs : Journaux d'exécution et informations d'erreur.
Query results : Les données de sortie de la tâche de débogage.
Étape 4 : Démarrer le nœud Flink SQL Streaming
-
Déployez le nœud Flink SQL Streaming.
La tâche doit être déployée dans Operation Center avant de pouvoir s'exécuter. Suivez les instructions à l'écran pour déployer le nœud. Pour plus d'informations, consultez Déployer un nœud.
RemarqueCette opération déploie également la tâche dans l'espace de travail Flink VVP. Vous pouvez afficher les tâches déployées via DataWorks dans Flink VVP Operation Center > Job O&M.
-
Démarrez le nœud Flink SQL Streaming.
Une fois la tâche déployée, cliquez sur Enable sous Deploy to production environment . Dans Operation Center, accédez à , recherchez la tâche que vous souhaitez démarrer et cliquez sur Script Parameters dans la colonne Go to operation and maintenance pour démarrer la tâche en temps réel et afficher son état d'exécution.