Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Règles CEP dynamiques : format JSON et utilisation

Dernière mise à jour :Aug 09, 2026

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é par patternName.

  • 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é par patternName.

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 :

  1. Avoir collecté un coupon de lieu (étape facultative).

  2. Avoir ajouté des articles à leur panier trois fois ou plus.

  3. 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.

Étapes suivantes