Realtime Compute for Apache Flink prend en charge le développement de tâches d'agent d'IA événementielles et en flux continu, basées sur le framework open source Apache Flink Agents. Cette rubrique explique les concepts fondamentaux des agents Flink, en couvrant les exigences de version, un guide de démarrage rapide, le développement personnalisé, la soumission de tâches et la gestion des dépendances.
Présentation
Apache Flink Agents est un nouveau sous-projet de la communauté Apache Flink dédié à la création d'agents d'IA événementiels. Il s'appuie sur le moteur de streaming éprouvé de Flink pour doter les agents d'IA de capacités distribuées, avec état, tolérantes aux pannes et fonctionnant en flux continu, les rendant ainsi prêts pour les environnements de production d'entreprise.
Flink Agents intègre des modules natifs pour les appels aux grands modèles de langage (LLM), l'utilisation d'outils, la gestion de la mémoire, l'orchestration dynamique et l'observabilité. Grâce à ces fonctionnalités, vous pouvez rapidement construire des agents d'IA fiables, conçus pour la production, fonctionnant en continu et à grande échelle.
Flink Agents offre les avantages principaux suivants :
|
Fonctionnalité |
Description |
|
Coordination distribuée |
Le framework fournit des capacités essentielles telles que le routage des événements, le partitionnement, la cohérence de l'état distribué et la récupération automatique après défaillance, éliminant ainsi la nécessité de gérer plusieurs réplicas d'agents. |
|
Mémoire avec état |
Flink State gère la mémoire des agents, qui est automatiquement persistée, restaurée et redistribuée lors des points de contrôle. Cela supprime le besoin d'un système de stockage externe pour maintenir la cohérence de l'état. |
|
Traitement en flux continu |
Les agents répondent en temps réel à mesure que les événements arrivent continuellement. Le moteur de traitement en flux continu de Flink constitue le cœur haute performance et à faible latence de ce modèle opérationnel. |
|
Fiabilité en production |
Les points de contrôle distribués permettent une récupération automatique en cas de défaillance des nœuds avec des garanties de traitement exactly-once, assurant qu'aucune donnée ni aucun événement n'est perdu. |
Types d'agents
Créez des agents Workflow ou ReAct selon vos cas d'utilisation.
|
Type |
Cas d'utilisation |
Description |
|
Agent Workflow |
Scénarios avec des processus bien définis |
Orchestre le comportement de l'agent via un workflow prédéfini et événementiel. Pour plus d'informations, consultez Workflow Agent. |
|
Agent ReAct |
Scénarios nécessitant une prise de décision flexible |
Combine le raisonnement et l'action, permettant à un LLM de déterminer de manière autonome les étapes d'exécution. Pour plus d'informations, consultez ReAct Agent. |
Versions et modèles
-
Lorsque vous soumettez une tâche Flink Agents, sélectionnez la version du moteur VVR VVR-11,8.preview.2 ou ultérieure. Sinon, la tâche échouera en raison de dépendances d'exécution manquantes.
Version du moteur VVR
Version Flink
Version Flink Agents
vvr-11,8.preview.2-jdk11-flink-1,20
1,20
0,2.1
Les modèles sont fournis par la plateforme. Vous n'avez pas besoin de demander une clé API distincte pour les appeler dans votre tâche. Cette fonctionnalité est en version bêta et nécessite un accès via une liste d'autorisation. Pour plus d'informations, consultez Service Flink AI (modèles intégrés).
Prérequis
-
Produits et autorisations
Activez Realtime Compute for Apache Flink et créez un espace de travail. Pour plus d'informations, consultez Activer Realtime Compute for Apache Flink.
Si vous utilisez un utilisateur RAM ou un rôle RAM, assurez-vous de disposer des autorisations nécessaires pour la console Flink. Pour plus d'informations, consultez Gestion des autorisations.
-
Environnement de développement local
Python : Python 3.10 ou 3.11.
Java : Java 11+ et Maven 3+.
Démarrage rapide
Cette section utilise un exemple d'analyse d'avis produits pour vous guider tout au long du processus. Vous allez créer un exemple de tâche Flink avec un agent d'analyse d'avis en utilisant le framework Flink Agents. Lorsque la tâche s'exécute, l'agent appelle le service de grand modèle de langage (LLM) de Model Studio pour analyser les avis produits, générant ainsi un score de satisfaction (de 1 à 5) et les raisons d'éventuelles insatisfactions.
Étape 1 : Télécharger les fichiers d'exemple de la tâche
Tâche Python
Téléchargez quickstart-python.zip, qui contient les fichiers suivants :
|
Fichier |
Description |
|
main.py |
Point d'entrée de la tâche. Il crée l'environnement d'exécution, enregistre la connexion LLM et construit le pipeline de traitement en flux continu. |
|
review_analysis_agent.py |
Implémentation de l'agent. Il définit le modèle d'invite, la configuration ChatModel et la logique de gestion des actions. |
Téléversez directement le package ZIP sans le décompresser. Les dépendances Python utilisées dans cet exemple, telles que openai et dashscope, sont préinstallées dans l'image du moteur VVR.
Tâche Java
Téléchargez les fichiers suivants :
|
Fichier |
Description |
|
Un JAR fat de test que vous pouvez téléverser directement. |
|
|
Le code source Java, fourni à titre de référence. |
Étape 2 : Téléverser et déployer la tâche
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la colonne Actions de l'espace de travail cible, cliquez sur Console.
Dans le volet de navigation de gauche, cliquez sur File Management. Puis, cliquez sur Upload Resource pour téléverser les fichiers de tâche que vous avez téléchargés.
Sur la page Operations > Jobs, cliquez sur Deploy Job. Sélectionnez Python Job ou JAR Job selon le type de votre tâche et renseignez les informations de déploiement.
Configuration de la tâche Python
|
Paramètre |
Exemple |
|
Mode de déploiement |
Mode Streaming |
|
Nom du déploiement |
flink-agents-quickstart-python |
|
Version du moteur |
vvr-11,8.preview.2-jdk11-flink-1,20 |
|
Adresse du fichier Python |
quickstart-python.zip |
|
Module d'entrée |
quickstart.main |
|
Cible de déploiement |
default-queue |
Configuration de la tâche Java
|
Paramètre |
Exemple |
|
Mode de déploiement |
Mode Streaming |
|
Nom du déploiement |
flink-agents-quickstart-java |
|
Version du moteur |
vvr-11,8.preview.2-jdk11-flink-1,20 |
|
URI JAR |
quickstart-java.jar |
|
Classe du point d'entrée |
org.apache.flink.agents.quickstart.Main |
|
Cible de déploiement |
default-queue |
Paramètres d'exécution
Accédez à Deployment Details > Runtime parameter configuration > Edit > Other Configurations, puis ajoutez les éléments de configuration Flink.
-
Tâche Python
python.executable: python3.10 python.client.executable: python3.10 containerized.master.env.FLINK_HOME: /flink containerized.taskmanager.env.FLINK_HOME: /flinkDescription des paramètres
Paramètre
Description
python.executable/python.client.executableSpécifie la version de l'interpréteur Python (pour les tâches Python uniquement).
containerized.{master,taskmanager}.env.FLINK_HOMEDéfinit la variable d'environnement
FLINK_HOMEpour le JobManager et les TaskManagers (pour les tâches Python uniquement).RemarqueConfigurez la variable d'environnement avec les préfixes
containerized.master.env.etcontainerized.taskmanager.env.pour garantir sa disponibilité à la fois pour le JobManager et les TaskManagers.
Étape 3 : Démarrer la tâche et afficher les résultats
Sur la page Operations > Jobs, localisez votre tâche et cliquez sur Start dans la colonne Actions.
Dans la boîte de dialogue Start Job, sélectionnez stateless start et cliquez sur Start.
L'exemple utilise une source de données en mémoire, donc la tâche se termine automatiquement une fois le traitement terminé. Attendez que le statut passe à finished, puis cliquez sur le nom de la tâche pour accéder à la page des détails.
Sous l'onglet TaskManager, consultez les journaux et recherchez le mot-clé
Review analysis result:pour voir la sortie.
Dans un environnement de production, vous utiliseriez généralement une source de données en flux continu comme Apache Kafka, et la tâche s'exécuterait en continu.
Développement d'un agent personnalisé
Installation de Flink Agents
Python
Nous vous recommandons d'utiliser un environnement virtuel :
python3 -m venv flink-agents-env
source flink-agents-env/bin/activate
pip install flink-agents
Pour plus d'informations, consultez le Guide d'installation de Flink Agents.
Java
Ajoutez la dépendance suivante à votre fichier pom.xml :
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-agents-api</artifactId>
<version>${flink-agents.version}</version>
<scope>provided</scope>
</dependency>
Déclarez les dépendances Flink Agents avec <scope>provided</scope> car elles sont fournies par le moteur VVR.
Pour déboguer dans votre IDE local, ajoutez la dépendance suivante :
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-agents-ide-support</artifactId>
<version>${flink-agents.version}</version>
<scope>provided</scope>
</dependency>
Écriture de l'agent
La structure centrale d'une tâche Flink Agents est la suivante :
Enregistrer une connexion LLM : Enregistrez une connexion ChatModel en utilisant
AgentsExecutionEnvironment. Configurez les informations de connexion, telles que la clé API et l'URL de l'endpoint, via unResourceDescriptor. Pour plus d'informations, consultez Chat Models.Définir l'agent : Implémentez votre agent personnalisé, configurez son modèle d'invite et définissez la logique de gestion des événements en utilisant l'annotation
@action(Python) ou@Action(Java).Construire le pipeline : Connectez le flux de données d'entrée à l'agent pour le traitement et produisez les résultats d'analyse.
Pour plus d'informations sur les concepts tels que l'agent, l'invite, l'outil et la mémoire, consultez la documentation de développement Flink Agents.
Tests locaux
Un MiniCluster Flink gère automatiquement les tests locaux ; vous n'avez donc pas besoin de déployer un cluster séparé.
Python
python your_agent_job.py
Java
mvn exec:java -Dexec.mainClass="com.example.YourAgentJob"
Pour plus d'informations, consultez la documentation de déploiement Flink Agents.
Gestion des dépendances
Dépendances préinstallées
L'image du moteur VVR inclut les bibliothèques principales de Flink Agents et certaines dépendances associées préinstallées.
Dépendances Python
|
Dépendance |
Description |
|
Bibliothèque principale flink-agents |
L'API Python de Flink Agents |
|
openai |
Interface compatible OpenAI (inclut la plateforme Model Studio) |
|
dashscope |
Alibaba Cloud Model Studio (interface native pour Qwen) |
|
mcp |
Model Context Protocol |
Téléversez les dépendances non préinstallées avec votre tâche. Pour plus d'informations, consultez Gérer les dépendances Python.
Dépendances Java
L'image du moteur VVR inclut Flink Agents et ses intégrations LLM intégrées. Empaquetez uniquement les dépendances tierces non intégrées dans le JAR de votre tâche.
Gestion des dépendances des tâches Java
Utilisez le plugin maven-shade-plugin pour empaqueter votre tâche dans un JAR fat. Configurez la portée des dépendances selon les règles suivantes :
Dépendances principales de Flink Agents (telles que
flink-agents-apietflink-agents-runtime) : Définissez la portée sur<scope>provided</scope>et ne les incluez pas dans le JAR fat.Dépendances principales de Flink (telles que
flink-streaming-java) : Définissez la portée sur<scope>provided</scope>.-
Dépendances d'intégration LLM : Traitez-les selon les cas suivants.
Pour les intégrations déjà prises en charge par la communauté Flink Agents (telles que OpenAI et Anthropic) : Définissez la portée sur
<scope>provided</scope>et ne les incluez pas dans le JAR fat. Pour la liste complète des intégrations prises en charge dans Flink Agents 0,2.1, consultez Built-in Providers.Pour les intégrations non encore prises en charge par la communauté : Définissez la portée sur
<scope>compile</scope>et incluez-les dans le JAR fat.
Bibliothèques en conflit avec Flink (telles que Jackson) : Configurez la relocation dans le plugin
maven-shade-plugin.
Pour un exemple complet de configuration pom.xml, consultez le fichier quickstart-java-src.zip.
Soumission de tâches à Realtime Compute for Flink
Tâche Python
Le processus de soumission est identique à celui d'une tâche PyFlink standard. Pour plus d'informations, consultez Développer une tâche Python. Les paramètres clés sont les suivants :
|
Paramètre |
Description |
|
Version du moteur |
Sélectionnez une version du moteur VVR qui prend en charge Flink Agents. Pour les versions spécifiques, consultez la section « Versions et modèles ». |
|
Adresse du fichier Python |
Le fichier Python d'entrée de la tâche ou le package ZIP. |
|
Module d'entrée |
Le nom du module d'entrée pour le package ZIP. |
|
Bibliothèques Python |
Dépendances supplémentaires. |
Tâche Java
Le processus de soumission est identique à celui d'une tâche JAR Flink standard. Pour plus d'informations, consultez Développer une tâche JAR. Les paramètres clés sont les suivants :
|
Paramètre |
Description |
|
Version du moteur |
Sélectionnez une version du moteur VVR qui prend en charge Flink Agents. Pour les versions spécifiques, consultez la section « Versions et modèles ». |
|
URI JAR |
Le chemin d'accès au JAR fat. |
|
Classe du point d'entrée |
Le nom qualifié complet de la classe d'entrée de la tâche. |
Paramètres d'exécution courants
Version de l'interpréteur Python
L'image du moteur VVR utilise par défaut Python 3.9, mais Flink Agents nécessite Python 3.10 ou 3.11. Configurez les éléments suivants pour une tâche Python :
python.executable: python3.10
python.client.executable: python3.10
containerized.master.env.FLINK_HOME: /flink
containerized.taskmanager.env.FLINK_HOME: /flink
Variables d'environnement personnalisées
Si le code de votre tâche lit la configuration à partir de variables d'environnement (telles qu'une URL d'endpoint ou une sélection de modèle), définissez les variables pour les processus JobManager et TaskManager en utilisant leurs préfixes respectifs :
containerized.master.env.<ENV_VAR_NAME>: <value>
containerized.taskmanager.env.<ENV_VAR_NAME>: <value>
Exemple de changement de modèle :
containerized.master.env.OPENAI_MODEL: qwen3.6-flash
containerized.taskmanager.env.OPENAI_MODEL: qwen3.6-flash