Le traitement complexe d'événements (CEP) Flink dynamique utilise un format JSON pour décrire les règles. Cette approche permet de stocker et de mettre à jour les règles CEP sans modifier ni recompiler le code Java.
Cette rubrique s'adresse aux :
Développeurs de plateformes de contrôle des risques familiarisés avec le CEP Flink dynamique qui souhaitent comprendre le schéma JSON et décider s'il convient d'ajouter des couches d'abstraction supplémentaires.
Responsables de la stratégie de contrôle des risques qui maîtrisent la logique métier mais ne possèdent pas de compétences en Java, et qui souhaitent rédiger ou ajuster directement les règles CEP en JSON.
Définition du format JSON
Une règle CEP se modélise sous forme de graphe orienté. Chaque nœud du graphe représente un motif d'événement. Chaque arête définit la stratégie de sélection des événements, c'est-à-dire la condition de transition d'un motif correspondant vers le suivant. Les graphes peuvent être imbriqués : un nœud de graphe peut devenir l'enfant d'un graphe plus vaste, permettant ainsi de regrouper des motifs.
Les sections suivantes décrivent chaque composant du schéma JSON.
Définition d'un nœud
Un nœud représente un motif unique et complet.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
name |
string | Oui | Nom unique du nœud. Les noms de nœuds doivent être uniques dans l'ensemble du graphe. |
type |
enum(string) | Oui | ATOMIC pour un nœud sans motif enfant. COMPOSITE pour un nœud contenant un motif enfant. |
quantifier |
dict | Oui | Décrit la manière dont le motif doit correspondre. Consultez la section Définition du quantificateur. |
condition |
dict | Non | Filtre les événements auxquels le motif s'applique. Consultez la section Définition de la condition. |
Définition du quantificateur
Un quantificateur décrit le nombre de fois où les événements doivent correspondre à un motif, ainsi que la stratégie de contiguïté au sein de ce motif. Par exemple, le motif A* possède une valeur properties définie sur LOOPING et une consumingStrategy définie sur SKIP_TILL_ANY.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
consumingStrategy |
enum(string) | Oui | Stratégie de sélection des événements au sein du motif. Valeurs valides : STRICT, SKIP_TILL_NEXT, SKIP_TILL_ANY. Consultez la section Définition de la contiguïté. |
times |
dict | Non | Nombre de fois où le motif doit correspondre. Voir l'exemple ci-dessous. |
properties |
array of enum(string) | Oui | Indicateurs de comportement de correspondance. Consultez la section Valeurs des propriétés du quantificateur. |
untilCondition |
dict | Non | Condition d'arrêt. Valide uniquement après un motif doté d'un quantificateur LOOPING. Consultez la section Définition de la condition. |
Exemple de valeur times :
"times": {
"from": 3,
"to": 3,
"windowTime": {
"unit": "MINUTES",
"size": 12
}
}
Les champs from et to sont des entiers. Le champ unit dans windowTime accepte les valeurs DAYS, HOURS, MINUTES, SECONDS ou MILLISECONDS. Définissez windowTime sur null pour omettre toute contrainte de temps par correspondance.
Définition de la condition
Une condition filtre les événements qui répondent à des critères spécifiques. Par exemple, « navigation pendant plus de 5 minutes » constitue une condition qui filtre les clients selon la durée de leur session.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
type |
enum(string) | Oui | Type de condition. Valeurs valides : CLASS, AVIATOR, GROOVY. |
| Champs personnalisés supplémentaires | — | Non | Tous les champs sérialisables supplémentaires spécifiques au type de condition. |
Quand utiliser chaque type de condition
| Scénario | Type recommandé |
|---|---|
| Logique métier nécessitant toute l'expressivité de Java ou une évaluation avec état sur les événements précédents | CLASS |
Comparaisons de seuils changeant fréquemment (par ex. price > 10) sans redéploiement du job |
AVIATOR |
| Logique multi-champs ou opérations sur les chaînes changeant fréquemment sans redéploiement du job | GROOVY |
Utilisez AVIATOR ou GROOVY lorsque vous devez mettre à jour les seuils de condition en modifiant une valeur dans la base de données, sans aucune modification ni recompilation du code requise.
Condition CLASS
Une condition CLASS délègue l'évaluation à une classe Java que vous fournissez.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
type |
enum(string) | Oui | Valeur fixe : CLASS. |
className |
string | Oui | Nom qualifié complet de la classe, par exemple com.alibaba.ververica.cep.demo.StartCondition. |
Condition CLASS avec paramètres personnalisés (CustomArgsCondition)
Une condition CLASS standard ne reçoit que le nom de la classe ; elle n'accepte pas de paramètres d'exécution. CustomArgsCondition étend la condition CLASS avec un tableau de chaînes (args) que le framework transmet lors de la construction de l'instance de condition. Cette approche permet de mettre à jour les paramètres de condition dans la base de données sans modifier ni recompiler la classe Java.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
type |
enum(string) | Oui | Valeur fixe : CLASS. |
className |
string | Oui | Nom qualifié complet de la classe, par exemple com.alibaba.ververica.cep.demo.CustomMiddleCondition. |
args |
array of string | Oui | Paramètres transmis au constructeur de la condition au moment de l'exécution. |
Condition d'expression Aviator
Aviator est un moteur d'évaluation d'expressions qui compile les expressions en bytecode au moment de l'exécution. Pour plus d'informations, consultez aviatorscript.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
type |
string | Oui | Valeur fixe : AVIATOR. |
expression |
string | Oui | Chaîne d'expression Aviator, telle que price > 10. Les variables de l'expression (par exemple, price) correspondent aux champs définis dans la classe d'événement Java. Mettez à jour cette chaîne dans la base de données pour modifier dynamiquement le seuil ; le job Flink CEP charge la nouvelle expression et crée une nouvelle AviatorCondition pour les événements suivants. |
Condition d'expression Groovy
Groovy est un langage typé dynamiquement pour la machine virtuelle Java (JVM). Pour plus d'informations sur la syntaxe Groovy, consultez Syntaxe Groovy.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
type |
string | Oui | Valeur fixe : GROOVY. |
expression |
string | Oui | Chaîne d'expression Groovy, telle que price > 5.0 && name.contains("mid"). Les variables correspondent aux champs de la classe d'événement Java. Mettez à jour cette chaîne dans la base de données pour modifier dynamiquement la logique ; le job Flink CEP charge la nouvelle chaîne Groovy et crée une nouvelle GroovyCondition pour les événements suivants. |
Définition d'une arête
Une arête relie deux nœuds de motif et définit la stratégie de sélection des événements pour cette transition.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
source |
string | Oui | Nom du nœud de motif source. |
target |
string | Oui | Nom du nœud de motif cible. |
type |
enum(string) | Oui | Stratégie de sélection des événements. Valeurs valides : STRICT, SKIP_TILL_NEXT, SKIP_TILL_ANY, NOT_FOLLOW, NOT_NEXT. Consultez la section Définition de la contiguïté. |
Définition de GraphNode
Un GraphNode représente une séquence de motifs complète. Il étend le nœud de base avec des champs de structure de graphe (nodes et edges) et des champs de politique (window et afterMatchSkipStrategy). Comme GraphNode est traité comme un sous-type de Node, un GraphNode peut être imbriqué dans un autre GraphNode pour créer des motifs groupés (GroupPattern).
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
name |
string | Oui | Nom unique du graphe. Les noms de graphes doivent être uniques. |
type |
enum(string) | Oui | Valeur fixe : COMPOSITE. |
version |
int | Oui | Version du format JSON. Valeur par défaut : 1. |
nodes |
array of Node | Oui | Motifs enfants dans ce graphe. Ne doit pas être vide. |
edges |
array of Edge | Oui | Connexions entre les motifs enfants. Peut être vide. |
window |
dict | Non | Contrainte de fenêtre temporelle. Voir la description ci-dessous. |
afterMatchSkipStrategy |
dict | Oui | Stratégie de saut appliquée après une correspondance complète. Consultez la section Définition de la stratégie de saut après correspondance. |
quantifier |
dict | Oui | Décrit la manière dont le motif de graphe global doit correspondre. Consultez la section Définition du quantificateur. |
Champ Window :
Le champ window contraint le temps autorisé pour une correspondance complète. Le champ type contrôle l'application de la limite de temps :
FIRST_AND_LAST: temps maximal entre le premier et le dernier événement d'une correspondance complète.PREVIOUS_AND_CURRENT: temps maximal entre les correspondances de deux motifs enfants adjacents quelconques.
Exemple :
"window": {
"type": "FIRST_AND_LAST",
"time": {
"unit": "DAYS",
"size": 1
}
}
Le champ unit accepte les valeurs DAYS, HOURS, MINUTES, SECONDS ou MILLISECONDS. La valeur size est un entier long ou integer.
Définition de la stratégie de saut après correspondance
La stratégie de saut après correspondance contrôle les correspondances partielles à rejeter une fois qu'une correspondance complète a été trouvée.
| Champ | Type | Obligatoire | Description |
|---|---|---|---|
type |
enum(string) | Oui | Stratégie de saut. Valeurs valides : NO_SKIP, SKIP_TO_NEXT, SKIP_PAST_LAST_EVENT, SKIP_TO_FIRST, SKIP_TO_LAST. |
patternName |
string | Non | Nom du motif utilisé par SKIP_TO_FIRST et SKIP_TO_LAST. |
Les stratégies ont les comportements suivants :
NO_SKIP(par défaut) : chaque correspondance réussie est émise sans rejet.SKIP_TO_NEXT: rejette chaque correspondance partielle qui a commencé avec le même événement que la correspondance actuelle.SKIP_PAST_LAST_EVENT: rejette chaque correspondance partielle qui a commencé entre le début et la fin de la correspondance actuelle.SKIP_TO_FIRST: rejette chaque correspondance partielle qui a commencé entre le début de la correspondance actuelle et la première occurrence de l'événement nommé parpatternName.SKIP_TO_LAST: rejette chaque correspondance partielle qui a commencé entre le début de la correspondance actuelle et la dernière occurrence de l'événement nommé parpatternName.
Pour plus d'informations, consultez Stratégie de saut après correspondance.
Définition de la contiguïté
La contiguïté contrôle la rigueur avec laquelle les événements doivent se suivre au sein d'un motif ou le long d'une arête.
| Valeur | Signification |
|---|---|
STRICT |
Contiguïté stricte. Aucun événement non correspondant ne peut apparaître entre les événements correspondants. |
SKIP_TILL_NEXT |
Contiguïté relâchée. Les événements non correspondants entre les événements correspondants sont ignorés silencieusement. |
SKIP_TILL_ANY |
Contiguïté relâchée non déterministe. Plus permissive que SKIP_TILL_NEXT — permet plusieurs correspondances pour certains événements correspondants. |
NOT_NEXT |
L'événement suivant immédiatement la source ne doit pas correspondre au motif cible. |
NOT_FOLLOW |
Aucun événement correspondant au motif cible ne peut apparaître n'importe où après la source. |
Pour plus d'informations, consultez FlinkCEP — Traitement complexe d'événements pour Flink.
Valeurs des propriétés du quantificateur
Les propriétés du quantificateur décrivent la cardinalité et la stratégie de correspondance d'un motif.
| Valeur | Signification |
|---|---|
SINGLE |
Le motif doit correspondre exactement une fois. |
LOOPING |
Le motif peut correspondre plusieurs fois en boucle, similaire à * et + dans les expressions régulières. |
TIMES |
Le motif doit correspondre un nombre spécifié de fois, tel que défini dans le champ times. |
GREEDY |
Lors de la correspondance, la séquence la plus longue possible est privilégiée. |
OPTIONAL |
Le motif est facultatif et peut ne pas correspondre du tout. |
Exemple 1 : Motif courant
Cet exemple utilise le CEP Flink dynamique pour identifier les clients qui doivent recevoir des offres marketing ajustées lors d'une promotion e-commerce en temps réel. Dans une fenêtre de 10 minutes, les clients ciblés doivent :
Avoir collecté un coupon de lieu (étape facultative).
Avoir ajouté des articles à leur panier trois fois ou plus.
Ne pas avoir finalisé le paiement.
Les trois conditions sont modélisées comme StartCondition, MiddleCondition et EndCondition. Le motif Java équivalent est :
Pattern<Event, Event> pattern =
Pattern.<Event>begin("start")
.where(new StartCondition())
.optional()
.followedBy("middle")
.where(new MiddleCondition())
.timesOrMore(3)
.notFollowedBy("end")
.where(new EndCondition())
.within(Time.minutes(10));
La règle JSON équivalente est :
{
"name": "end",
"quantifier": {
"consumingStrategy": "SKIP_TILL_NEXT",
"properties": [
"SINGLE"
],
"times": null,
"untilCondition": null
},
"condition": null,
"nodes": [
{
"name": "end",
"quantifier": {
"consumingStrategy": "SKIP_TILL_NEXT",
"properties": [
"SINGLE"
],
"times": null,
"untilCondition": null
},
"condition": {
"className": "com.alibaba.ververica.cep.demo.condition.EndCondition",
"type": "CLASS"
},
"type": "ATOMIC"
},
{
"name": "middle",
"quantifier": {
"consumingStrategy": "SKIP_TILL_NEXT",
"properties": [
"LOOPING"
],
"times": {
"from": 3,
"to": 3,
"windowTime": null
},
"untilCondition": null
},
"condition": {
"className": "com.alibaba.ververica.cep.demo.condition.MiddleCondition",
"type": "CLASS"
},
"type": "ATOMIC"
},
{
"name": "start",
"quantifier": {
"consumingStrategy": "SKIP_TILL_NEXT",
"properties": [
"SINGLE",
"OPTIONAL"
],
"times": null,
"untilCondition": null
},
"condition": {
"className": "com.alibaba.ververica.cep.demo.condition.StartCondition",
"type": "CLASS"
},
"type": "ATOMIC"
}
],
"edges": [
{
"source": "middle",
"target": "end",
"type": "NOT_FOLLOW"
},
{
"source": "start",
"target": "middle",
"type": "SKIP_TILL_NEXT"
}
],
"window": {
"type": "FIRST_AND_LAST",
"time": {
"unit": "MINUTES",
"size": 10
}
},
"afterMatchStrategy": {
"type": "NO_SKIP",
"patternName": null
},
"type": "COMPOSITE",
"version": 1
}
Exemple 2 : Condition avec paramètres personnalisés
Cet exemple montre comment appliquer différentes stratégies marketing à des clients de classes différentes lors d'un événement promotionnel e-commerce en temps réel. Par exemple, vous pouvez envoyer des messages texte marketing aux clients de classe A, envoyer des coupons aux clients de classe B et ne prendre aucune action marketing pour les autres clients.
Avec une condition CLASS standard, la classe est codée en dur pour gérer les niveaux A et B. Si vous souhaitez ajouter la classe C ou ajuster la stratégie, vous devez réécrire et recompiler le code de déploiement. Pour simplifier cette opération, utilisez une condition avec des paramètres personnalisés (CustomArgsCondition). Après avoir défini dans le code la manière d'ajuster les stratégies en fonction du paramètre transmis, il vous suffit de modifier la valeur du paramètre args dans la base de données, sans aucune modification ni recompilation du code requise.
L'extrait de code suivant montre la condition initialement définie dans le motif :
"condition": {
"args": [
"A", "B"
],
"className": "org.apache.flink.cep.pattern.conditions.CustomMiddleCondition",
"type": "CLASS"
}
Pour ajouter le niveau C à la stratégie, mettez à jour le tableau args dans la base de données :
"condition": {
"args": [
"A", "B", "C"
],
"className": "org.apache.flink.cep.pattern.conditions.CustomMiddleCondition",
"type": "CLASS"
}
Pour un exemple complet et fonctionnel, consultez Démo.
aviatorscript et Démo sont des ressources tierces. Ces liens peuvent être lents à charger ou temporairement indisponibles.