Cette rubrique explique comment déployer et démarrer des tâches PyFlink en streaming et par lots, en détaillant le flux de développement dans Realtime Compute for Apache Flink.
Prérequis
Si vous utilisez un utilisateur RAM ou un rôle RAM pour accéder à la console, assurez-vous que l'identité dispose des autorisations requises. Pour plus d'informations, consultez la section Gestion des autorisations.
Un espace de travail a été créé. Pour plus d'informations, consultez la section Activer Realtime Compute for Apache Flink.
Étape 1 : Préparer les fichiers de code Python
La console de gestion de Realtime Compute for Apache Flink ne fournit pas d'environnement de développement Python. Développez vos tâches localement. Pour plus d'informations sur le débogage des tâches et les connecteurs, consultez la section Développer des tâches PyFlink.
Assurez-vous que la version de Flink utilisée pour le développement local correspond à la version du moteur sélectionnée à l'Étape 3 : Déployer une tâche PyFlink. Pour savoir comment utiliser d'autres dépendances, telles que les environnements virtuels Python personnalisés, les packages Python tiers, les packages JAR et les fichiers de données, consultez la section Utiliser les dépendances Python.
Pour vous aider à démarrer rapidement, cette rubrique fournit des exemples de fichiers Python pour une tâche de comptage de mots ainsi qu'un fichier de données exemple. Vous pouvez les télécharger et les utiliser dans les étapes suivantes.
-
Téléchargez le fichier de tâche Python exemple approprié.
Tâche en streaming : word_count_streaming.py.
Tâche par lot : word_count_batch.py.
Cliquez sur Shakespeare pour télécharger le fichier de données exemple.
Étape 2 : Télécharger les fichiers Python et de données
Connectez-vous à la console Realtime Compute.
Recherchez l'espace de travail Flink cible et cliquez sur Console dans la colonne Actions.
Dans le volet de navigation de gauche, cliquez sur Artifacts.
-
Cliquez sur Upload Artifact pour télécharger les fichiers Python et de données.
Téléchargez les exemples de fichiers Python et de données obtenus à l'étape 1. Pour plus d'informations sur les chemins de stockage des fichiers, consultez la section Artifacts.
Étape 3 : Déployer une tâche PyFlink
Streaming
Sur la page , cliquez sur .
-
Configurez les paramètres de déploiement.
Paramètre
Description
Exemple
Deployment mode
Sélectionnez le mode streaming.
stream mode
Deployment name
Saisissez un nom pour le déploiement Python.
flink-streaming-test-python
Engine version
Version du moteur Flink pour le déploiement.
Nous vous recommandons d'utiliser une version portant le tag RECOMMENDED ou STABLE pour une meilleure fiabilité et de meilleures performances. Pour plus d'informations, consultez les Notes de version et les Versions du moteur.
vvr-8.0.9-flink-1.17
Python URI
Téléchargez le fichier exemple word_count_streaming.py. Cliquez ensuite sur l'icône de téléchargement
pour sélectionner et envoyer le fichier.Si le fichier existe déjà dans les Artifacts, vous pouvez le sélectionner directement sans le retélécharger.
-
Entry module
Module de point d'entrée du programme.
-
Ce paramètre n'est pas requis si la tâche PyFlink est un fichier .py.
-
Si la tâche PyFlink est un fichier .zip, vous devez saisir le module de point d'entrée. Exemple :
word_count.
Non requis
Entry point main arguments
Arguments à transmettre à la méthode principale.
Pour ce tutoriel, saisissez le chemin de stockage du fichier de données d'entrée, Shakespeare.
--input oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/ShakespeareVous pouvez copier le chemin complet du fichier Shakespeare depuis la page Artifacts.
Deployment target
Dans la liste déroulante, sélectionnez la Resource Queue ou le session cluster cible. Les clusters de session ne sont pas recommandés pour la production. Pour plus d'informations, consultez les sections Gérer les files d'attente et Créer un cluster de session.
ImportantLes déploiements sur un cluster de session ne prennent pas en charge les métriques de surveillance, la configuration des alertes ni Autopilot. Utilisez les clusters de session uniquement pour le développement et les tests ; ne les utilisez pas dans les environnements de production. Pour plus d'informations, consultez la section Déboguer les déploiements.
default-queue
Pour plus d'informations sur les autres paramètres de configuration, consultez la section Déployer une tâche.
-
Cliquez sur Deploy.
Batch
Sur la page , cliquez sur Create Deployment et sélectionnez Python Deployment.
-
Configurez les paramètres de déploiement.
Paramètre
Description
Exemple
Deployment mode
Sélectionnez le mode batch.
batch mode
Deployment name
Saisissez un nom pour le déploiement Python.
flink-batch-test-python
Engine version
Version du moteur Flink pour le déploiement.
Nous vous recommandons d'utiliser une version portant le tag RECOMMENDED ou STABLE pour une meilleure fiabilité et de meilleures performances. Pour plus d'informations, consultez les Notes de version et les Versions du moteur.
vvr-8.0.9-flink-1.17
Python URI
Téléchargez le fichier exemple word_count_batch.py. Cliquez ensuite sur l'icône de téléchargement
pour sélectionner et envoyer le fichier.-
Entry module
Module de point d'entrée du programme.
-
Ce paramètre n'est pas requis si la tâche PyFlink est un fichier .py.
-
Si la tâche PyFlink est un fichier .zip, vous devez saisir le module de point d'entrée. Exemple :
word_count.
Non requis
Entry point main arguments
Arguments à transmettre à la méthode principale.
Pour ce tutoriel, saisissez les chemins de stockage du fichier d'entrée Shakespeare et du répertoire de sortie
python-batch-quickstart-test-output.RemarqueVous devez uniquement spécifier le chemin du répertoire de sortie. Le répertoire de sortie doit se trouver dans le même répertoire parent que le fichier d'entrée. Vous n'avez pas besoin de créer le répertoire de sortie au préalable.
--input oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/Shakespeare--output oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/python-batch-quickstart-test-outputVous pouvez copier le chemin complet du fichier Shakespeare depuis la page Artifacts.
Deployment target
Dans la liste déroulante, sélectionnez la Resource Queue ou le session cluster cible. Les clusters de session ne sont pas recommandés pour la production. Pour plus d'informations, consultez les sections Gérer les files d'attente et Créer un cluster de session.
ImportantLes déploiements sur un cluster de session ne prennent pas en charge les métriques de surveillance, la configuration des alertes ni Autopilot. Utilisez les clusters de session uniquement pour le développement et les tests ; ne les utilisez pas dans les environnements de production. Pour plus d'informations, consultez la section Déboguer les déploiements.
default-queue
Pour plus d'informations sur les autres paramètres de configuration, consultez la section Déployer une tâche.
-
Cliquez sur Deploy.
Étape 4 : Démarrer le déploiement et afficher les résultats
Streaming
Sur la page , recherchez le déploiement cible et cliquez sur Start dans la colonne Actions.
-
Dans la boîte de dialogue Start Job, sélectionnez Initial Mode et cliquez sur Start. Pour plus d'informations, consultez la section Démarrer un déploiement.
Après avoir cliqué sur Start, un statut RUNNING ou FINISHED indique que le déploiement s'exécute comme prévu. Si vous utilisez le fichier exemple de cette rubrique, le statut final est FINISHED.
-
Une fois que le statut du déploiement passe à RUNNING, affichez les résultats du déploiement en streaming.
ImportantSi vous utilisez le fichier Python exemple de cette rubrique, les résultats sont supprimés lorsque le déploiement en streaming passe à l'état FINISHED. Par conséquent, vous ne pouvez afficher les résultats que lorsque le déploiement est à l'état RUNNING.
Dans le fichier journal TaskManager se terminant par .out, recherchez
shakespearepour trouver le résultat du calcul.Dans l'onglet Logs, cliquez sur l'onglet Running Task Managers. Pour le TaskManager concerné, cliquez sur le sous-onglet Log List. Ouvrez le fichier
flink.outet saisissezshakespearedans la zone de recherche en haut à droite pour localiser le résultat du comptage de mots, tel que(shakespeare,1).
Batch
-
Sur la page , recherchez le déploiement cible et cliquez sur Start dans la colonne Actions.
Pour filtrer la liste, sélectionnez Batch Deployment dans la liste déroulante de type.
Dans la boîte de dialogue Start Job. Pour plus d'informations, consultez la section Démarrer un déploiement.
-
Une fois que le statut du déploiement passe à Start, affichez les résultats du déploiement par lot.
Connectez-vous à la console OSS. Accédez au répertoire oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/python-batch-quickstart-test-output. Cliquez sur le dossier nommé d'après la date et l'heure de démarrage du déploiement, cliquez sur le fichier cible, puis cliquez sur Actions dans le panneau qui s'affiche.
Le déploiement par lot produit un fichier .ext. Après avoir téléchargé le fichier, ouvrez-le avec un éditeur de texte ou Microsoft Word pour afficher les résultats. La sortie ressemble à ce qui suit :
(As,40) (At,5) (Ay,1) (Be,9) (By,14) (Do,4) (He,7) (I,,4) (If,34) (In,36) (Is,10) (It,6)
(Facultatif) Étape 5 : Arrêter un déploiement
Pour appliquer des modifications à une tâche (telles que des modifications de code, des mises à jour des paramètres WITH ou des changements de version), vous devez la redéployer, l'arrêter, puis la redémarrer. Un redémarrage est également nécessaire pour un démarrage sans état ou pour appliquer des modifications de configuration non dynamiques. Pour plus d'informations sur l'arrêt d'une tâche, consultez la section Arrêter une tâche.
Rubriques connexes
Vous pouvez configurer les ressources d'un déploiement avant de le démarrer ou modifier les ressources une fois le déploiement en cours d'exécution. Deux modes de configuration des ressources sont pris en charge : basique (granularité grossière) et expert (granularité fine). Pour plus d'informations, consultez la section Configurer les ressources de déploiement.
Realtime Compute for Apache Flink prend en charge les mises à jour dynamiques des paramètres de déploiement. Cela permet aux configurations de prendre effet plus rapidement et réduit les temps d'arrêt du service causés par l'arrêt et le démarrage des déploiements. Pour plus d'informations, consultez la section Mise à l'échelle dynamique et mises à jour des paramètres.
Configurez les niveaux de journalisation du déploiement et spécifiez différentes sorties pour différents niveaux de journal. Pour plus d'informations, consultez la section Configurer la sortie des journaux.
Pour une présentation détaillée du flux de développement SQL, consultez la section Tâche Flink SQL.
Construire un entrepôt de données en temps réel avec Hologres.
Construire un lac de données en streaming avec Paimon et StarRocks.