O processamento complexo de eventos (CEP) dinâmico do Flink utiliza o formato JSON para descrever regras. Isso permite armazenar e atualizar regras CEP sem modificar ou recompilar código Java.
Este tópico destina-se a:
Desenvolvedores de plataformas de controle de risco com conhecimento em CEP dinâmico do Flink que desejam compreender o esquema JSON e avaliar a necessidade de camadas adicionais de abstração.
Profissionais de estratégia de controle de risco que dominam a lógica de controle de risco, mas não possuem experiência com Java e pretendem escrever ou ajustar regras CEP diretamente em JSON.
Definição do formato JSON
Uma regra CEP é modelada como um grafo direcionado. Cada nó representa um padrão de evento, e cada aresta define a estratégia de seleção de eventos — a condição de transição de um padrão correspondente para o próximo. Grafos podem ser aninhados: um nó pode ser filho de um grafo maior, o que possibilita padrões agrupados.
As seções a seguir detalham cada componente do esquema JSON.
Definição de nó
Um nó representa um padrão único e completo.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
string |
Sim |
Nome exclusivo do nó. Os nomes dos nós devem ser únicos em todo o grafo. |
|
|
enum(string) |
Sim |
|
|
|
dict |
Sim |
Descreve como corresponder ao padrão. Consulte Definição de quantificador. |
|
|
dict |
Não |
Filtra os eventos aos quais o padrão se aplica. Consulte Definição de condição. |
Definição de quantificador
Um quantificador especifica quantas vezes os eventos devem corresponder a um padrão e a estratégia de contiguidade dentro desse padrão. Por exemplo, o padrão A* tem o valor properties definido como LOOPING e uma consumingStrategy de SKIP_TILL_ANY.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
enum(string) |
Sim |
Estratégia de seleção de eventos dentro do padrão. Valores válidos: |
|
|
dict |
Não |
Quantidade de vezes que o padrão deve corresponder. Veja o exemplo abaixo. |
|
|
array of enum(string) |
Sim |
Sinalizadores de comportamento de correspondência. Consulte Valores de propriedade do quantificador. |
|
|
dict |
Não |
Condição de parada. Válida apenas após um padrão com quantificador |
Exemplo de valor times:
"times": {
"from": 3,
"to": 3,
"windowTime": {
"unit": "MINUTES",
"size": 12
}
}
Os campos from e to são inteiros. O campo unit em windowTime aceita DAYS, HOURS, MINUTES, SECONDS ou MILLISECONDS. Defina windowTime como null para omitir a restrição de tempo por correspondência.
Definição de condição
Uma condição filtra eventos que atendem a critérios específicos. Por exemplo, "navegação por mais de 5 minutos" é uma condição que filtra clientes pela duração da sessão.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
enum(string) |
Sim |
Tipo da condição. Valores válidos: |
|
Campos personalizados adicionais |
— |
Não |
Quaisquer campos serializáveis adicionais específicos do tipo de condição. |
Quando usar cada tipo de condição
|
Cenário |
Tipo recomendado |
|
Lógica de negócios que exige expressividade total do Java ou avaliação com estado em eventos anteriores |
|
|
Comparações de limiar que mudam frequentemente (por exemplo, |
|
|
Lógica multicampo ou operações de string que mudam frequentemente sem reimplantar o job |
|
Use AVIATOR ou GROOVY quando precisar atualizar limiares de condição alterando um valor no banco de dados, sem alterar ou recompilar código.
Condição CLASS
Uma condição CLASS delega a execução a uma classe Java fornecida por você.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
enum(string) |
Sim |
Valor fixo: |
|
|
string |
Sim |
Nome totalmente qualificado da classe, como |
Condição CLASS com parâmetros personalizados (CustomArgsCondition)
Uma condição CLASS padrão recebe apenas o nome da classe e não aceita parâmetros em tempo de execução. A CustomArgsCondition estende a condição CLASS com um array de strings (args) que o framework transmite ao construir a instância da condição. Esse recurso permite atualizar parâmetros de condição no banco de dados sem alterar ou recompilar a classe Java.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
enum(string) |
Sim |
Valor fixo: |
|
|
string |
Sim |
Nome totalmente qualificado da classe, como |
|
|
array of string |
Sim |
Parâmetros transmitidos ao construtor da condição em tempo de execução. |
Condição de expressão Aviator
Aviator é um mecanismo de avaliação de expressões que compila expressões para bytecode em tempo de execução. Para mais informações, consulte aviatorscript.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
string |
Sim |
Valor fixo: |
|
|
string |
Sim |
String de expressão Aviator, como |
Condição de expressão Groovy
Groovy é uma linguagem de tipagem dinâmica para a Java Virtual Machine (JVM). Para mais informações sobre a sintaxe do Groovy, consulte Sintaxe do Groovy.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
string |
Sim |
Valor fixo: |
|
|
string |
Sim |
String de expressão Groovy, como |
Definição de aresta
Uma aresta conecta dois nós de padrão e define a estratégia de seleção de eventos para essa transição.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
string |
Sim |
Nome do nó de padrão de origem. |
|
|
string |
Sim |
Nome do nó de padrão de destino. |
|
|
enum(string) |
Sim |
Estratégia de seleção de eventos. Valores válidos: |
Definição de GraphNode
Um GraphNode representa uma sequência completa de padrões. Ele estende o nó básico com campos de estrutura de grafo (nodes e edges) e campos de política (window e afterMatchSkipStrategy). Como GraphNode é tratado como subtipo de Node, é possível aninhar um GraphNode dentro de outro GraphNode para criar padrões agrupados (GroupPattern).
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
string |
Sim |
Nome exclusivo do grafo. Os nomes dos grafos devem ser únicos. |
|
|
enum(string) |
Sim |
Valor fixo: |
|
|
int |
Sim |
Versão do formato JSON. Valor padrão: |
|
|
array of Node |
Sim |
Padrões filhos neste grafo. Não pode estar vazio. |
|
|
array of Edge |
Sim |
Conexões entre os padrões filhos. Pode estar vazio. |
|
|
dict |
Não |
Restrição de janela de tempo. Consulte a descrição abaixo. |
|
|
dict |
Sim |
Estratégia de salto aplicada após uma correspondência completa. Consulte Definição de estratégia de salto pós-correspondência. |
|
|
dict |
Sim |
Descreve como corresponder ao padrão geral do grafo. Consulte Definição de quantificador. |
Campo Window:
O campo window restringe o tempo permitido para uma correspondência completa. O type controla a aplicação do limite de tempo:
FIRST_AND_LAST: tempo máximo entre o primeiro e o último evento em uma correspondência completa.PREVIOUS_AND_CURRENT: tempo máximo entre correspondências de quaisquer dois padrões filhos adjacentes.
Exemplo:
"window": {
"type": "FIRST_AND_LAST",
"time": {
"unit": "DAYS",
"size": 1
}
}
O campo unit aceita DAYS, HOURS, MINUTES, SECONDS ou MILLISECONDS. O valor size é um long ou integer.
Definição de estratégia de salto pós-correspondência
A estratégia de salto pós-correspondência controla quais correspondências parciais descartar após encontrar uma correspondência completa.
|
Campo |
Tipo |
Obrigatório |
Descrição |
|
|
enum(string) |
Sim |
Estratégia de salto. Valores válidos: |
|
|
string |
Não |
Nome do padrão usado por |
As estratégias apresentam os seguintes comportamentos:
NO_SKIP(padrão): Emite todas as correspondências bem-sucedidas sem descartar nenhuma.SKIP_TO_NEXT: Descarta todas as correspondências parciais iniciadas com o mesmo evento da correspondência atual.SKIP_PAST_LAST_EVENT: Descarta todas as correspondências parciais iniciadas entre o início e o fim da correspondência atual.SKIP_TO_FIRST: Descarta todas as correspondências parciais iniciadas entre o início da correspondência atual e a primeira ocorrência do evento nomeado porpatternName.SKIP_TO_LAST: Descarta todas as correspondências parciais iniciadas entre o início da correspondência atual e a última ocorrência do evento nomeado porpatternName.
Para mais informações, consulte Estratégia de salto pós-correspondência.
Definição de contiguidade
A contiguidade controla o rigor com que os eventos devem se suceder dentro de um padrão ou ao longo de uma aresta.
|
Valor |
Significado |
|
|
Contiguidade estrita. Nenhum evento não correspondido pode aparecer entre eventos correspondidos. |
|
|
Contiguidade relaxada. Eventos não correspondidos entre eventos correspondidos são ignorados silenciosamente. |
|
|
Contiguidade relaxada não determinística. Mais permissiva que |
|
|
O evento imediatamente seguinte à origem não deve corresponder ao padrão de destino. |
|
|
Nenhum evento correspondente ao padrão de destino pode aparecer após a origem. |
Para mais informações, consulte FlinkCEP — Processamento complexo de eventos para Flink.
Valores de propriedade do quantificador
As propriedades do quantificador descrevem a cardinalidade e a estratégia de correspondência para um padrão.
|
Valor |
Significado |
|
|
O padrão deve corresponder exatamente uma vez. |
|
|
O padrão pode corresponder várias vezes em loop, semelhante a |
|
|
O padrão deve corresponder a um número específico de vezes, conforme definido no campo |
|
|
Durante a correspondência, prioriza a sequência mais longa possível. |
|
|
O padrão é opcional e pode não corresponder. |
Exemplo 1: Padrão comum
Este exemplo usa o CEP dinâmico do Flink para identificar clientes elegíveis para receber ofertas de marketing ajustadas durante uma promoção de e-commerce em tempo real. Dentro de uma janela de 10 minutos, os clientes-alvo devem ter:
Coletado um cupom do local (etapa opcional).
Adicionado itens ao carrinho três ou mais vezes.
Não concluído o pagamento.
As três condições são modeladas como StartCondition, MiddleCondition e EndCondition. O padrão Java equivalente é:
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));
A regra JSON equivalente é:
{
"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
}
Exemplo 2: Condição com parâmetros personalizados
Este exemplo demonstra como aplicar diferentes estratégias de marketing a clientes de classes distintas durante um evento promocional de e-commerce em tempo real. Por exemplo, você pode enviar mensagens de texto de marketing para clientes da Classe A, enviar cupons para clientes da Classe B e não executar nenhuma ação de marketing para outros clientes.
Em uma condição CLASS padrão, a classe é codificada rigidamente para lidar com as categorias A e B. Para adicionar a Classe C ou ajustar a estratégia, é necessário reescrever e recompilar o código de implantação. Para simplificar, use uma condição com parâmetros personalizados (CustomArgsCondition). Após definir no código como ajustar as estratégias com base no parâmetro transmitido, basta alterar o valor do parâmetro args no banco de dados, sem alterar ou recompilar código.
O código de exemplo a seguir mostra a condição inicialmente definida no padrão:
"condition": {
"args": [
"A", "B"
],
"className": "org.apache.flink.cep.pattern.conditions.CustomMiddleCondition",
"type": "CLASS"
}
Para adicionar a categoria C à estratégia, atualize o array args no banco de dados:
"condition": {
"args": [
"A", "B", "C"
],
"className": "org.apache.flink.cep.pattern.conditions.CustomMiddleCondition",
"type": "CLASS"
}
Para um exemplo funcional completo, consulte Demo.
aviatorscript e Demo são recursos de terceiros. Esses links podem apresentar carregamento lento ou indisponibilidade temporária.