Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Flink CEP dinâmico

Última atualização: Jun 27, 2026

O Realtime Compute for Apache Flink oferece suporte a tarefas do Flink CEP com regras atualizadas dinamicamente em jobs DataStream. Este tópico usa um cenário de marketing em tempo real para demonstrar como criar um job do Flink CEP que carrega regras dinamicamente para processar dados de um tópico Kafka upstream.

Casos de uso

Graças à sua arquitetura distribuída, latência na casa dos milissegundos e recursos avançados de expressão de regras, o Flink CEP é ideal para diversas aplicações. Três cenários típicos incluem:

  • Controle de risco em tempo real: o Flink CEP identifica usuários em situação de risco. Por exemplo, é possível analisar logs de comportamento do cliente para sinalizar usuários que realizam mais de 10 transferências, totalizando mais de 10.000, em até 5 minutos.

  • Marketing em tempo real: o Flink CEP otimiza estratégias de marketing. Ao analisar logs de comportamento do usuário durante uma promoção de e-commerce, você identifica usuários que adicionam mais de três itens ao carrinho em 10 minutos, mas não concluem a compra, permitindo ajustes direcionados de marketing. O Flink CEP também é eficaz em cenários antifraude para marketing em tempo real.

  • Internet das Coisas (IoT): o Flink CEP detecta estados anormais e envia alertas. Por exemplo, ele emite um alerta de risco se uma bicicleta compartilhada sair de uma área designada e não retornar em até 15 minutos. Também pode ser combinado com sensores de IoT para detectar anomalias em linhas de produção. Se um sensor de temperatura relatar continuamente temperaturas acima de um limiar definido por três períodos consecutivos, um alerta será acionado.

Exemplo prático

Este tópico demonstra como usar o CEP dinâmico para atender a esses cenários. Neste exemplo, os logs de comportamento do cliente são armazenados no ApsaraMQ for Kafka. O job do Flink CEP consome esses dados enquanto consulta periodicamente uma tabela de regras em um banco de dados ApsaraDB RDS for MySQL. Ele busca as regras mais recentes adicionadas por um administrador de políticas e as utiliza para corresponder eventos. Quando ocorre uma correspondência, o job envia um alerta ou grava as informações relevantes em outro armazenamento de dados. A figura a seguir mostra o pipeline de dados geral.

Esta demonstração inicia primeiro o job do Flink CEP e depois insere a Regra 1, que corresponde a uma sequência onde três eventos consecutivos com uma action igual a 0 são seguidos por um evento onde a ação não é 1. Isso significa que um usuário visualizou um produto três vezes sem efetuar a compra.

Pré-requisitos

Procedimento

Este tópico descreve como escrever e atualizar dinamicamente um job do Flink CEP que monitora e registra usuários cujos logs de comportamento correspondem a regras específicas.

Etapa 1: Preparar dados de teste

Preparar o tópico Kafka upstream

  1. Faça login no console do ApsaraMQ for Kafka.

  2. Crie um tópico chamado demo_topic para armazenar logs simulados de comportamento do usuário.

    Para mais informações, consulte Etapa 1: Criar um tópico.

Preparar o banco de dados RDS

No console do Data Management (DMS), prepare os dados de teste para o ApsaraDB RDS for MySQL.

  1. Faça login na instância do ApsaraDB RDS for MySQL com uma conta privilegiada.

    Para mais informações, consulte Fazer login em uma instância do ApsaraDB RDS for MySQL usando o DMS.

  2. Crie a tabela de regras rds_demo para armazenar as regras do job do Flink CEP. Crie a tabela match_results para armazenar os dados que correspondem às regras.

    Na janela ativa do SQLConsole, insira os seguintes comandos e clique em Execute.

    CREATE DATABASE cep_demo_db;
    USE cep_demo_db;
    CREATE TABLE rds_demo (
      `id` VARCHAR(64),
      `version` INT,
      `pattern` VARCHAR(4096),
      `function` VARCHAR(512)
    );
    CREATE TABLE match_results (
        rule_id INT,
        rule_version INT,
        user_id INT,
        user_name VARCHAR(255),
        production_id INT,
        PRIMARY KEY (rule_id,rule_version,user_id,production_id)
    );

    Cada linha na tabela de regras rds_demo representa uma única regra. Ela inclui um id e uma version para distinguir entre diferentes regras e suas versões, um campo pattern que descreve o objeto pattern da API CEP e um campo function que descreve como processar a sequência de eventos correspondente ao padrão.

    Cada linha na tabela match_results representa uma correspondência onde o comportamento de um usuário para um produto específico está em conformidade com uma determinada regra. Esse registro pode ser usado para formular estratégias de vendas correspondentes, como enviar cupons para produtos relacionados.

Etapa 2: Configurar lista de permissões de IP

Para permitir que o job do Flink acesse a instância do ApsaraDB RDS for MySQL, adicione o bloco CIDR do workspace do Realtime Compute for Apache Flink à lista de permissões de endereços IP da instância.

  1. Obtenha o bloco CIDR da VPC do workspace do Realtime Compute for Apache Flink.

    1. Faça login no console do Realtime Compute for Apache Flink.

    2. Na coluna Actions do workspace desejado, escolha More > Workspace details.

    3. Na caixa de diálogo Workspace details, visualize o CIDR block do vSwitch do Flink totalmente gerenciado.

  2. Adicione o bloco CIDR do Flink totalmente gerenciado à lista de permissões de endereços IP da sua instância do ApsaraDB RDS for MySQL.

    Para mais informações, consulte Configurar uma lista de permissões de endereços IP. Na caixa de diálogo Modify whitelist, insira o bloco CIDR do Flink na caixa de texto IP addresses in whitelist. Se você tiver vários blocos CIDR, separe-os por vírgulas. Em seguida, clique em OK.

Etapa 3: Desenvolver e iniciar o job CEP

Nota

Todo o código deste tópico está disponível em nosso repositório GitHub. Para fins de demonstração, o código de exemplo neste tópico foi ligeiramente modificado na branch timeOrMoreAndWindow. Você pode baixar o arquivo completo ververica-cep-demo-master.zip para referência.

  1. Adicione flink-cep como dependência do projeto no arquivo POM Maven do job.

    Para mais informações sobre como lidar com outros pacotes JAR relacionados ao Flink e resolver conflitos, consulte Configurar dependências de ambiente do Flink.

    <dependency>
        <groupId>com.alibaba.ververica</groupId>
        <artifactId>flink-cep</artifactId>
        <version>1.17-vvr-8.0.8</version>
        <scope>provided</scope>
    </dependency>
  2. Desenvolva o código do job.

    1. Construa um Kafka Source.

      Para detalhes sobre como escrever o código, consulte Kafka DataStream Connector.

    2. Construa a API CEP.dynamicPatterns().

      Para suportar alterações dinâmicas de regras e correspondência de múltiplas regras para CEP, o Realtime Compute for Apache Flink define a API CEP.dynamicPatterns(). A API é definida da seguinte forma.

      public static <T, R> SingleOutputStreamOperator<R> dynamicPatterns(
               DataStream<T> input,
               PatternProcessorDiscovererFactory<T> discovererFactory,
               TimeBehaviour timeBehaviour,
               TypeInformation<R> outTypeInfo)

      A tabela a seguir descreve os parâmetros desta API. Você pode atualizar os valores dos parâmetros com base no seu caso de uso real.

      Parâmetro

      Descrição

      DataStream<T> input

      O fluxo de eventos de entrada.

      PatternProcessorDiscovererFactory<T> discovererFactory

      Uma fábrica que cria um PatternProcessorDiscoverer. Este descobridor busca as regras mais recentes e constrói as instâncias correspondentes de PatternProcessor.

      TimeBehaviour timeBehaviour

      Descreve como o job do Flink CEP lida com atributos de tempo dos eventos. Valores válidos:

      • TimeBehaviour.ProcessingTime: Processa eventos com base no tempo de processamento.

      • TimeBehaviour.EventTime: Processa eventos com base no tempo do evento.

      TypeInformation<R> outTypeInfo

      Descreve as informações de tipo do fluxo de saída.

      Para mais informações sobre conceitos comuns do Flink, como DataStream, TimeBehaviour e TypeInformation, consulte API DataStream, Tempo de Evento e Tempo de Processamento e TypeInformation.

      A interface PatternProcessor é um componente chave. Um PatternProcessor contém um Pattern específico que descreve como corresponder eventos e uma PatternProcessFunction que descreve como lidar com uma correspondência, como enviar um alerta. Ele também inclui um id e uma version para identificar o PatternProcessor. Para mais contexto, consulte a proposta.

      O patternProcessorDiscovererFactory cria um descobridor para buscar o PatternProcessor mais recente. O código de exemplo inclui uma classe abstrata que mostra como consultar periodicamente um armazenamento externo em busca de novas instâncias de PatternProcessor.

      public abstract class PeriodicPatternProcessorDiscoverer<T>
              implements PatternProcessorDiscoverer<T> {
          ...
          @Override
          public void discoverPatternProcessorUpdates(
                  PatternProcessorManager<T> patternProcessorManager) {
              // Periodically discovers the pattern processor updates.
              timer.schedule(
                      new TimerTask() {
                          @Override
                          public void run() {
                              if (arePatternProcessorsUpdated()) {
                                  List<PatternProcessor<T>> patternProcessors = null;
                                  try {
                                      patternProcessors = getLatestPatternProcessors();
                                  } catch (Exception e) {
                                      e.printStackTrace();
                                  }
                                  patternProcessorManager.onPatternProcessorsUpdated(patternProcessors);
                              }
                          }
                      },
                      0,
                      intervalMillis);
          }
          ...
      }

      O Realtime Compute for Apache Flink fornece uma implementação do JDBCPeriodicPatternProcessorDiscoverer para buscar as regras mais recentes de um banco de dados que suporta o protocolo JDBC, como ApsaraDB RDS for MySQL ou Hologres. Ao utilizá-lo, especifique os seguintes parâmetros.

      Parâmetro

      Descrição

      jdbcUrl

      A URL JDBC do banco de dados.

      jdbcDriver

      O nome da classe do driver do banco de dados.

      tableName

      O nome da tabela do banco de dados.

      initialPatternProcessors

      O PatternProcessor padrão a ser usado quando a tabela de regras no banco de dados estiver vazia.

      intervalMillis

      O intervalo de consulta ao banco de dados, em milissegundos.

      No seu código, use-o da seguinte maneira. O job imprimirá as regras correspondentes na saída do TaskManager do Flink.

      // import ......
      public class CepDemo {
          public static void main(String[] args) throws Exception {
              ......
              // DataStream Source
              DataStreamSource<Event> source =
                      env.fromSource(
                              kafkaSource,
                              WatermarkStrategy.<Event>forMonotonousTimestamps()
                                      .withTimestampAssigner((event, ts) -> event.getEventTime()),
                              "Kafka Source");
              env.setParallelism(1);
              // Key by userId and productionId.
              // Note: Only events with the same key will be processed to see if there is a match.
              KeyedStream<Event, Tuple2<Integer, Integer>> keyedStream =
                      source.assignTimestampsAndWatermarks(
                              WatermarkStrategy.<Event>forGenerator(ctx -> new EventBoundedOutOfOrdernessWatermarks(Duration.ofSeconds(5)))
                      ).keyBy(new KeySelector<Event, Tuple2<Integer, Integer>>() {
                          @Override
                          public Tuple2<Integer, Integer> getKey(Event value) throws Exception {
                              return Tuple2.of(value.getId(), value.getProductionId());
                          }
                      });
              SingleOutputStreamOperator<String> output =
                      CEP.dynamicPatterns(
                              keyedStream,
                              new JDBCPeriodicPatternProcessorDiscovererFactory<>(
                                      params.get(JDBC_URL_ARG),
                                      JDBC_DRIVE,
                                      params.get(TABLE_NAME_ARG),
                                      null,
                                      Long.parseLong(params.get(JDBC_INTERVAL_MILLIS_ARG))),
                              Boolean.parseBoolean(params.get(USING_EVENT_TIME)) ? TimeBehaviour.EventTime : TimeBehaviour.ProcessingTime,
                              TypeInformation.of(new TypeHint<String>() {}));
              output.print();
              // Compile and submit the job
              env.execute("CEPDemo");
          }
      }
      Nota

      Para fins de demonstração, o código de exemplo agrupa o fluxo de dados de entrada por id e product_id antes de conectá-lo ao CEP.dynamicPatterns(). Isso significa que apenas eventos com o mesmo id e product_id são considerados para correspondência de regras. Eventos com chaves diferentes não serão comparados entre si.

  3. No console do Realtime Compute for Apache Flink, faça upload do arquivo JAR e implante o job JAR. Para mais informações, consulte Implantar um job.

    Para começar rapidamente, baixe o arquivo de teste cep-demo.jar. A tabela a seguir descreve os parâmetros a serem configurados durante a implantação.

    Nota

    Como a fonte Kafka upstream está vazia e a tabela de regras do banco de dados não contém dados, o job não produzirá nenhuma saída após o início.

    Parâmetro

    Descrição

    Deployment mode

    Selecione Stream Mode.

    Deployment name

    Insira um nome para o job JAR.

    Engine version

    Para mais informações sobre versões do mecanismo, consulte Versões do mecanismo e Políticas de ciclo de vida. Recomendamos usar uma versão recomendada ou estável. As tags de versão são descritas da seguinte forma:

    • Versão recomendada: A versão secundária mais recente da versão principal atual.

    • Versão estável: A versão secundária mais recente de uma versão principal que ainda está dentro do período de suporte do produto e teve defeitos históricos corrigidos.

    • Versão normal: Outras versões secundárias que ainda estão dentro do período de suporte do produto.

    • Versão EOS: Uma versão que excedeu seu período de suporte do produto.

    JAR URL

    Faça upload do seu arquivo JAR empacotado ou do arquivo JAR de teste fornecido.

    Entry Point Class

    Insira com.alibaba.ververica.cep.demo.CepDemo.

    Entry Point Main Arguments

    Se você estiver usando seu próprio job desenvolvido e já tiver configurado as informações de armazenamento upstream e downstream, deixe este campo em branco. No entanto, se estiver usando o JAR de teste fornecido, configure este parâmetro. O código é o seguinte.

    --kafkaBrokers YOUR_KAFKA_BROKERS 
    --inputTopic YOUR_KAFKA_TOPIC 
    --inputTopicGroup YOUR_KAFKA_TOPIC_GROUP 
    --jdbcUrl jdbc:mysql://YOUR_DB_URL:port/DATABASE_NAME?user=YOUR_USERNAME&password=YOUR_PASSWORD
    --tableName YOUR_TABLE_NAME  
    --jdbcIntervalMs 3000
    --usingEventTime false

    A tabela a seguir descreve os parâmetros.

    • kafkaBrokers: O endereço do broker Kafka.

    • inputTopic: O nome do tópico Kafka.

    • inputTopicGroup: O grupo de consumidores Kafka.

    • jdbcUrl: A URL JDBC do banco de dados.

      Nota

      O nome de usuário e a senha na URL JDBC para este exemplo devem ser de uma conta padrão, e a senha pode conter apenas letras e números. Em um cenário real, você pode usar diferentes métodos de autenticação no seu job conforme necessário.

    • tableName: O nome da tabela de destino.

    • jdbcIntervalMs: O intervalo de consulta ao banco de dados.

    • usingEventTime: Especifica se deve usar o tempo do evento para processamento (true/false).

    Nota
    • Substitua os valores de espaço reservado pelas suas informações reais de armazenamento upstream e downstream.

    • Evite usar senhas em texto simples em um ambiente de produção. Recomendamos usar o recurso de gerenciamento de variáveis. Para mais informações, consulte Gerenciamento de variáveis.

  4. Na página Deployment details, na seção Other configuration, adicione os seguintes parâmetros de tempo de execução do job.

    Em aplicações práticas, o JAR flink-cep é carregado pelo carregador de classes do sistema, enquanto as classes relacionadas ao aviator são normalmente empacotadas em um JAR de usuário e carregadas pelo carregador de classes do usuário. Ao usar as duas configurações abaixo, você garante que o carregador de classes do sistema possa acessar classes no JAR do usuário ao tentar carregar uma classe, evitando assim falhas no carregamento de classes.

    kubernetes.application-mode.classpath.include-user-jar: 'true' 
    classloader.resolve-order: parent-first

    Para mais informações sobre como configurar parâmetros de tempo de execução, consulte Configuração de parâmetros de tempo de execução.

  5. Na página O&M > Deployments, localize a implantação desejada e clique em Start na coluna Actions.

    Para mais informações sobre como configurar parâmetros de inicialização do job, consulte Iniciar um job.

Etapa 4: Inserir uma regra

Com o job do Flink CEP em execução, insira a Regra 1: após três eventos consecutivos com uma action igual a 0, a ação do próximo evento não é 1. Isso significa que um usuário visualiza um produto três vezes sem efetuar a compra.

  1. Faça login no console do ApsaraDB RDS for MySQL.

  2. Insira a regra de atualização dinâmica.

    Concatene a string JSON com o id, a version e o nome da classe de função e, em seguida, insira-a no RDS.

    INSERT INTO rds_demo (
     `id`,
     `version`,
     `pattern`,
     `function`
    ) values(
      '1',
       1,
      '{"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":"start","quantifier":{"consumingStrategy":"SKIP_TILL_NEXT","properties":["LOOPING"],"times":{"from":3,"to":3,"windowTime":null},"untilCondition":null},"condition":{"expression":"action == 0","type":"AVIATOR"},"type":"ATOMIC"}],"edges":[{"source":"start","target":"end","type":"SKIP_TILL_NEXT"}],"window":null,"afterMatchStrategy":{"type":"SKIP_PAST_LAST_EVENT","patternName":null},"type":"COMPOSITE","version":1}',
      'com.alibaba.ververica.cep.demo.dynamic.DemoPatternProcessFunction')
    ;

    Para melhorar a usabilidade e a legibilidade do campo pattern no banco de dados, o Realtime Compute for Apache Flink define um formato de regra baseado em JSON. Para mais informações, consulte Formato JSON para regras no CEP dinâmico. O campo pattern na instrução SQL anterior contém uma string JSON serializada. Essa string representa um padrão que corresponde à seguinte sequência: após três eventos consecutivos com uma action igual a 0, a ação do próximo evento não é 1.

    Nota

    No código EndCondition, a condição definida é action != 1.

    • A descrição correspondente da API CEP é a seguinte.

      Pattern<Event, Event> pattern =
          Pattern.<Event>begin("start", AfterMatchSkipStrategy.skipPastLastEvent())
              .where(new StartCondition("action == 0"))
              .timesOrMore(3)
              .followedBy("end")
              .where(new EndCondition());
    • Converta-o para a string JSON correspondente usando o método em CepJsonUtils.

      public void printTestPattern(Pattern<?, ?> pattern) throws JsonProcessingException {
          System.out.println(CepJsonUtils.convertPatternToJSONString(pattern));
      }
    • A string JSON correspondente é a seguinte.

      Exemplo de string JSON de regra CEP dinâmica

      {
        "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": "start",
            "quantifier": {
              "consumingStrategy": "SKIP_TILL_NEXT",
              "properties": [
                "LOOPING"
              ],
              "times": {
                "from": 3,
                "to": 3,
                "windowTime": null
              },
              "untilCondition": null
            },
            "condition": {
              "expression": "action == 0",
              "type": "AVIATOR"
            },
            "type": "ATOMIC"
          }
        ],
        "edges": [
          {
            "source": "start",
            "target": "end",
            "type": "SKIP_TILL_NEXT"
          }
        ],
        "window": null,
        "afterMatchStrategy": {
          "type": "SKIP_PAST_LAST_EVENT",
          "patternName": null
        },
        "type": "COMPOSITE",
        "version": 1
      }
  3. Envie mensagens para o tópico demo_topic usando um cliente Kafka.

    Nesta demonstração, você também pode usar a página Start to send and consume message fornecida pelo ApsaraMQ for Kafka para enviar mensagens de teste.

    1,Ken,0,1,1662022777000
    1,Ken,0,1,1662022778000
    1,Ken,0,1,1662022779000
    1,Ken,0,1,1662022780000

    Selecione o método de envio Console, insira 1 no campo Message key, cole os dados de teste no campo Message content, defina Send to specified partition como No e envie a mensagem. A página exibirá uma notificação Message sent successfully.

    A tabela a seguir descreve os campos em demo_topic.

    Parâmetro

    Descrição

    id

    O ID do usuário.

    username

    O nome de usuário.

    action

    A ação do usuário. Valores válidos:

    • 0: operação de visualização

    • 1: ação de compra

    product_id

    O ID do produto.

    event_time

    O horário do evento de comportamento.

  4. Visualize a regra mais recente nos logs do JobManager e a correspondência nos logs do TaskManager.

    • Nos logs do JobManager, pesquise por JDBCPeriodicPatternProcessorDiscoverer para visualizar a regra mais recente.

      Para encontrar os logs, acesse Logs > JobManager e clique na aba Logs. Insira a palavra-chave na caixa de pesquisa para localizar a entrada de log relevante. A mensagem de log PatternProcessors have been updated confirma que a regra foi atualizada com sucesso.

    • No arquivo de log do TaskManager terminado em .out, pesquise por A match for Pattern of (id, version): (1, 1) para visualizar a correspondência resultante.

      Na página de detalhes da implantação, clique na aba Logs, selecione Running task managers e abra a subaba Logs para o TaskManager correspondente. Pesquise pela palavra-chave no arquivo flink.out para localizar a sequência de eventos correspondente.

  5. Consulte a tabela match_results executando SELECT * FROM para ver os resultados que correspondem à regra.

    A consulta retorna um registro com os campos rule_id, rule_version, user_id, user_name e production_id, com valores 1, 1, 1, Ken e 1, respectivamente.

Etapa 5: Atualizar regra de correspondência

Estratégias de marketing frequentemente têm restrições de tempo. Esta etapa atualiza a regra para exigir que os três eventos action = 0 ocorram dentro de um intervalo de 15 minutos.

  1. Defina o parâmetro usingEventTime como true.

    1. Na página O&M > Deployments, localize a implantação desejada e clique em Cancel na coluna Actions.

    2. Em Deployment Details > Entry Point Main Arguments , clique em Edit, defina o parâmetro usingEventTime como true e clique em Save.

    3. Clique em Start para iniciar o job novamente.

  2. Insira a nova regra.

    A descrição correspondente da API CEP é a seguinte.

    Pattern<Event, Event> pattern =
            Pattern.<Event>begin("start", AfterMatchSkipStrategy.skipPastLastEvent())
                    .where(new StartCondition("action == 0"))
                    .timesOrMore(3,Time.minutes(15))
                    .followedBy("end")
                    .where(new EndCondition());
    printTestPattern(pattern);

    Insira a nova regra na tabela rds_demo.

    # To avoid rule conflicts for this demo, delete the previous rule first.
    DELETE FROM `rds_demo` WHERE `id` = 1;
    # Insert the new rule: three consecutive `action = 0` events within 15 minutes, followed by an event where the action is not 1. The rule version is (1, 2).
    INSERT INTO rds_demo (`id`,`version`,`pattern`,`function`) values('1',2,'{"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":"start","quantifier":{"consumingStrategy":"SKIP_TILL_NEXT","properties":["LOOPING"],"times":{"from":3,"to":3,"windowTime":{"unit":"MINUTES","size":15}},"untilCondition":null},"condition":{"expression":"action == 0","type":"AVIATOR"},"type":"ATOMIC"}],"edges":[{"source":"start","target":"end","type":"SKIP_TILL_NEXT"}],"window":null,"afterMatchStrategy":{"type":"SKIP_PAST_LAST_EVENT","patternName":null},"type":"COMPOSITE","version":1}','com.alibaba.ververica.cep.demo.dynamic.DemoPatternProcessFunction');
  3. No console do Kafka, envie oito mensagens para acionar uma correspondência.

    A seguir estão oito mensagens de exemplo.

    2,Tom,0,1,1739584800000   #10:00
    2,Tom,0,1,1739585400000   #10:10
    2,Tom,0,1,1739585700000   #10:15
    2,Tom,0,1,1739586000000   #10:20
    3,Ali,0,1,1739586600000   #10:30
    3,Ali,0,1,1739588400000   #11:00
    3,Ali,0,1,1739589000000   #11:10
    3,Ali,0,1,1739590200000   #11:30
  4. Consulte a tabela match_results executando SELECT * FROM para ver os resultados que correspondem à regra.

    A consulta retorna dois registros, com colunas para rule_id, rule_version, user_id, user_name e production_id. O primeiro registro é para Ken (rule_version=1) e o segundo é para Tom (rule_version=2), ambos associados ao production_id=1.

    Os resultados mostram que apenas o comportamento de Tom corresponde à nova regra, pois as ações de Ali ocorreram ao longo de mais de 15 minutos. Para promoções por tempo limitado, isso permite enviar cupons para usuários que visitam repetidamente um produto dentro de um período específico, incentivando a compra.