Este tópico descreve as instruções de processamento de eventos complexas (CEP) compatíveis com o Realtime Compute for Apache Flink.
Informações básicas
Comparadas ao CEP SQL do Apache Flink, as instruções CEP do Realtime Compute for Apache Flink oferecem recursos aprimorados, como saída de eventos correspondentes que não chegam dentro de um intervalo de tempo específico, contiguidade relaxada por meio de followedBy() e configuração do padrão de contiguidade entre eventos. Para obter mais informações sobre os recursos básicos do CEP SQL do Apache Flink, consulte Pattern Recognition.
Limites
Apenas o Realtime Compute for Apache Flink com versão de mecanismo vvr-6.0.2-flink-1.15 ou posterior é compatível com a sintaxe estendida de CEP SQL.
Apenas o Realtime Compute for Apache Flink com versão de mecanismo vvr-6.0.5-flink-1.15 ou posterior é compatível com padrões de grupo e a sintaxe AFTER MATCH NO SKIP.
Saída de eventos correspondentes que não chegam dentro de um intervalo de tempo específico
O exemplo a seguir mostra uma sequência de eventos de entrada.
+----+------+------------------+
| id | type | rowtime |
+----+------+------------------+
| 1 | A | 2022-09-19 12:00 |
| 2 | B | 2022-09-19 12:01 |
| 3 | A | 2022-09-19 12:02 |
| 4 | B | 2022-09-19 12:05 |
+----+------+------------------+
Para especificar que o intervalo de tempo entre os eventos no padrão A B deve ser de no máximo 2 minutos, adicione WITHIN INTERVAL '2' MINUTES após a instrução PATTERN. Exemplo de código:
SELECT *
FROM MyTable MATCH_RECOGNIZE (
ORDER BY rowtime
MEASURES
A.id AS aid,
B.id AS bid,
A.rowtime AS atime,
B.rowtime AS btime
PATTERN (A B) WITHIN INTERVAL '2' MINUTES
DEFINE
A AS type = 'A',
B AS type = 'B'
) AS T
Sem a cláusula WITHIN, o sistema obtém duas sequências correspondentes: id=1, id=2 e id=3, id=4. Com a cláusula WITHIN, apenas a primeira sequência correspondente é retornada. Isso ocorre porque o intervalo de tempo entre o Evento A e o Evento B na segunda sequência é de 3 minutos, excedendo o limite de 2 minutos definido na cláusula WITHIN. A seguinte saída é retornada:
+-----+-----+------------------+------------------+
| aid | bid | atime | btime |
+-----+-----+------------------+------------------+
| 1 | 2 | 2022-09-19 12:00 | 2022-09-19 12:01 |
+-----+-----+------------------+------------------+
Com a inclusão da cláusula WITHIN, as sequências em que os eventos não chegam dentro do intervalo especificado são consideradas falhas de correspondência e descartadas. Para obter essas sequências mesmo quando os eventos ultrapassam o tempo limite, utilize a instrução ONE ROW PER MATCH SHOW TIMEOUT MATCHES. Exemplo de código:
SELECT *
FROM MyTable MATCH_RECOGNIZE (
ORDER BY rowtime
MEASURES
A.id AS aid,
B.id AS bid,
A.rowtime AS atime,
B.rowtime AS btime
ONE ROW PER MATCH SHOW TIMEOUT MATCHES
PATTERN (A B) WITHIN INTERVAL '2' MINUTES
DEFINE
A AS type = 'A',
B AS type = 'B'
) AS T
A saída abaixo inclui a sequência de eventos correspondentes em que os eventos não chegaram dentro do intervalo de tempo especificado.
+-----+--------+------------------+------------------+
| aid | bid | atime | btime |
+-----+--------+------------------+------------------+
| 1 | 2 | 2022-09-19 12:00 | 2022-09-19 12:01 |
| 3 | <NULL> | 2022-09-19 12:00 | <NULL> |
+-----+--------+------------------+------------------+
O momento de chegada do Evento B com ID 4 excede o intervalo definido na cláusula WITHIN. Como esse evento não faz parte das sequências correspondentes válidas, os campos bid e btime associados a ele são nulos.
Padrões de contiguidade entre eventos
A e a API Java CEP do Apache Flink aceitam os seguintes padrões de contiguidade entre eventos: contiguidade estrita com next(), contiguidade relaxada com followedBy(), contiguidade relaxada não determinística com followedByAny(), não contiguidade estrita com notNext() e não contiguidade relaxada com notFollowedBy().
Por padrão, o CEP SQL do Apache Flink utiliza contiguidade estrita. Nesse padrão, todos os eventos correspondentes devem aparecer rigorosamente um após o outro, sem eventos não correspondentes entre eles. No exemplo anterior, o padrão (A B) exige que o Evento A e o Evento B ocorram em sequência estrita. O Realtime Compute for Apache Flink estende essa funcionalidade para oferecer capacidade de expressão totalmente equivalente à API Java.
A tabela a seguir descreve as sequências de correspondência para diferentes padrões, considerando a sequência de eventos de entrada a1, b1, a2, a3, b2, b3.
Durante o processo de correspondência, a cláusula AFTER MATCH SKIP utiliza a estratégia SKIP TO NEXT ROW. Para mais detalhes sobre as estratégias da cláusula AFTER MATCH SKIP, consulte ou After Match Strategy.
|
API Java |
SQL |
Estratégia |
Sequência correspondente |
|
|
|
Contiguidade estrita: exige que todos os eventos correspondentes apareçam rigorosamente um após o outro, sem permitir eventos não correspondentes entre eles. |
|
|
|
C é um caractere indefinido na cláusula DEFINE e indica qualquer correspondência. |
Contiguidade relaxada: ignora eventos não correspondentes que aparecem entre os eventos correspondentes. |
|
|
|
C é um caractere indefinido na cláusula DEFINE e indica qualquer correspondência. |
Contiguidade relaxada não determinística: flexibiliza ainda mais a contiguidade e permite correspondências adicionais que ignoram eventos específicos. |
Nota
As sequências correspondentes neste exemplo foram obtidas usando a estratégia SKIP TO NEXT ROW, padrão de AFTER MATCH no CEP SQL. Já a API Java CEP do Apache Flink usa NO_SKIP como padrão. Para saber como utilizar a estratégia AFTER MATCH NO SKIP, consulte a seção Estratégia AFTER MATCH NO SKIP deste tópico. |
|
|
|
Não contiguidade estrita: espera-se que nenhum evento correspondente apareça logo após um evento correspondente. |
|
|
|
C é um caractere indefinido na cláusula DEFINE e indica qualquer correspondência. Nota
Para usar |
Não contiguidade relaxada: espera-se que um evento correspondente não apareça entre dois eventos correspondentes. Quando usada com a cláusula WITHIN no final do padrão, garante que nenhum evento de um tipo específico ocorra dentro de um período determinado. |
Sem correspondência |
Contiguidade e correspondência gulosa em padrões de repetição
O CEP SQL não aceita contiguidade relaxada não determinística em padrões de repetição.
A e a API Java CEP do Apache Flink permitem definir estratégias de contiguidade e correspondência gulosa dentro de padrões de repetição. Por padrão, o CEP SQL do Apache Flink aplica contiguidade estrita e correspondência gulosa. Por exemplo, no padrão A+, nenhum outro evento é permitido entre múltiplos Eventos A, e o sistema tenta corresponder ao maior número possível de Eventos A. Adicione um ou mais pontos de interrogação (?) após o quantificador de repetição, como *, + ou {3, }, para ajustar as estratégias de contiguidade e correspondência gulosa.
A tabela abaixo detalha as sequências de correspondência para diferentes padrões, dada a sequência de entrada a1, b1, a2, a3, c1 e a condição A AS type = 'a', C AS type = 'a' or type = 'c'.
Durante o processo de correspondência, a cláusula AFTER MATCH SKIP utiliza a estratégia SKIP TO NEXT ROW. Para mais detalhes sobre as estratégias da cláusula AFTER MATCH SKIP, consulte ou After Match Strategy.
|
Identificadores |
Contiguidade |
Estratégia de correspondência gulosa |
Padrão de exemplo |
Semântica equivalente |
Sequência correspondente |
|
Nenhum |
Contiguidade estrita |
Gulosa |
|
|
|
|
? |
Contiguidade estrita |
Não gulosa |
|
|
|
|
?? |
Contiguidade relaxada |
Gulosa |
|
|
|
|
??? |
Contiguidade relaxada |
Não gulosa |
|
|
|
until(condition) em padrões de repetição
A e a API Java CEP do Apache Flink permitem usar a função until(condition) para definir uma condição de parada em um padrão de repetição. Se o evento atual atender à condição especificada pela função until(condition), a correspondência do padrão de repetição termina imediatamente e a correspondência do próximo padrão começa a partir desse evento. Em implantações SQL do Realtime Compute for Apache Flink, anexe a sintaxe { CONDITION } a um quantificador de repetição, como +, * ou {3, }, para expressar a semântica de until.
A tabela a seguir apresenta as sequências de correspondência para diferentes padrões, considerando a sequência de entrada a1, d1, a2, b1, a3, c1 e a condição DEFINE A AS A.type = 'a' OR A.type = 'b', B AS B.type = 'b', C AS C.type = 'c'.
Durante o processo de correspondência, a cláusula AFTER MATCH SKIP utiliza a estratégia SKIP TO NEXT ROW. Para mais detalhes sobre as estratégias da cláusula AFTER MATCH SKIP, consulte ou After Match Strategy.
|
Padrão |
Semântica equivalente |
Sequência correspondente |
Descrição |
|
|
|
|
Eventos iniciados com a ou b podem corresponder ao padrão de repetição A, aplicando-se contiguidade estrita entre os eventos desse padrão e os padrões A e C. Como existe d1 entre a1 e a2 na sequência de entrada, a correspondência não pode começar com a1. |
|
|
|
|
A condição |
|
|
|
|
Aplica-se contiguidade relaxada entre os padrões A e C. O padrão de repetição iniciado em a2 termina em b1, pulando b1 e a3 para corresponder a c1. |
|
|
|
|
A contiguidade relaxada aplica-se aos eventos no padrão de repetição A. O padrão pula d1 e termina em b1 para corresponder a a1 e a2. |
Padrão de grupo
A e a API Java CEP do Apache Flink aceitam padrões de grupo. Nesses padrões, múltiplos subpadrões são combinados e usados nas funções next(), followedBy() ou followedByAny(). Um padrão de grupo pode ser repetido como um todo. Em implantações SQL do Realtime Compute for Apache Flink, utilize a sintaxe padrão SQL (...) para definir um padrão de grupo. Quantificadores de repetição como +, * e {3, } são permitidos.
Por exemplo, no padrão PATTERN (A (B C*)+? D) , (B C*) representa um padrão de grupo declarado para aparecer mais de uma vez. O ponto de interrogação (?) indica o uso da estratégia de correspondência não gulosa. Exemplo de código Java:
Pattern.<String>begin("A").where(...)
.next(
Pattern.<String>begin("B").where(...)
.next("C").where(...).oneOrMore().optional().greedy().consecutive())
.oneOrMore().consecutive()
.next("D").where(...)
A cláusula MEASURES define o conteúdo da saída de um padrão de grupo correspondente. Suponha que as sequências correspondentes sejam b1, b2 c1 e b3 c2 c3 sempre que o padrão de grupo especificado for encontrado. Nesse caso, use a cláusula MEASURES para exibir apenas parte dos resultados. Se precisar mostrar apenas o Evento b na saída, utilize FIRST(B.id) para obter o Evento b da primeira sequência correspondente e FIRST(B.id,1) para o Evento b da segunda sequência. Aplique o mesmo método para a terceira sequência. Assim, a saída do padrão de grupo será b1, b2, b3. Exemplo de código:
SELECT *
FROM MyTable MATCH_RECOGNIZE (
ORDER BY rowtime
MEASURES
FIRST(B.id) AS b1_id,
FIRST(B.id,1) AS b2_id,
FIRST(B.id,2) AS b3_id
PATTERN (A (B C*)+? D)
DEFINE
A AS type = 'A',
B AS type = 'B',
C AS type = 'C',
D AS type = 'D'
) AS T
Observe que a contiguidade declarada entre um padrão de grupo e seu padrão anterior aplica-se apenas ao primeiro padrão dentro do grupo, e não ao grupo inteiro. Por exemplo, no padrão PATTERN (A {- X*? -} (B C)) , o uso de followedBy entre o padrão A e o padrão de grupo (B C) estabelece uma relação de followedBy entre o padrão A e o padrão B. Isso significa que vários eventos que não correspondem ao padrão B podem existir entre o padrão A e o padrão de grupo (B C), mas eventos que não correspondem ao padrão de grupo (B C) não são permitidos. Se o padrão PATTERN (A {- X*? -} (B C)) não gerar saída para a sequência de entrada a1 b1 d1 b2 c1, isso ocorre porque o processo de correspondência entra imediatamente no padrão de grupo (B C) após o surgimento de b1, e d1 falha ao tentar corresponder ao padrão C. Como resultado, a correspondência da sequência falha.
Padrões de grupo com repetição, como
PATTERN ((A B)+), não aceitam correspondência gulosa.Não é possível usar padrões de grupo, como
PATTERN (A+{(B C)})ePATTERN (A [^(B C)]), nas sintaxesuntilounotNext.O primeiro padrão em um padrão de grupo, como em
PATTERN (A (B? C)), não pode ser declarado como opcional.
Estratégia AFTER MATCH NO SKIP
Na e na API Java CEP do Apache Flink, a estratégia padrão de AFTER MATCH é NO_SKIP. Já no CEP SQL do Apache Flink, a estratégia padrão é SKIP_TO_NEXT_ROW. O Realtime Compute for Apache Flink estende a cláusula AFTER MATCH do padrão SQL, permitindo declarar a estratégia NO_SKIP por meio da cláusula AFTER MATCH NO SKIP. Com a estratégia NO_SKIP, os processos de correspondência existentes não terminam nem são descartados após a conclusão da correspondência de uma sequência.
Geralmente, a estratégia NO_SKIP é usada em conjunto com followedByAny para pular eventos específicos em cenários de contiguidade relaxada. Por exemplo, dada a sequência de entrada a1 b1 b2 b3 c1, a saída para o padrão PATTERN (A {- X* -} B {- Y*? -} C) será a1 b1 c1 ao usar a estratégia padrão AFTER MATCH SKIP TO NEXT ROW. Esse padrão equivale a Pattern.begin("A").followedByAny("B").followedBy("C"). Isso acontece porque todas as sequências iniciadas com a1 são descartadas assim que a correspondência de a1 b1 c1 é concluída. No entanto, ao utilizar AFTER MATCH NO SKIP, todas as sequências correspondentes são obtidas. Nesse cenário, a1 b1 c1, a1 b2 c1 e a1 b3 c1 são retornados.