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
Um workspace do Realtime Compute for Apache Flink criado. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
Se você utilizar um usuário RAM ou uma função RAM, garanta que as permissões necessárias foram concedidas para o console do Flink. Para mais informações, consulte gerenciamento de permissões.
-
Armazenamento upstream e downstream:
Uma instância do ApsaraDB RDS for MySQL criada. Para mais informações, consulte Criar uma instância do ApsaraDB RDS for MySQL.
Uma instância do ApsaraMQ for Kafka criada. Para mais informações, consulte ApsaraMQ for Kafka.
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
Faça login no console do ApsaraMQ for Kafka.
-
Crie um tópico chamado
demo_topicpara 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.
-
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.
-
Crie a tabela de regras
rds_demopara armazenar as regras do job do Flink CEP. Crie a tabelamatch_resultspara 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_demorepresenta uma única regra. Ela inclui umide umaversionpara distinguir entre diferentes regras e suas versões, um campopatternque descreve o objeto pattern da API CEP e um campofunctionque descreve como processar a sequência de eventos correspondente ao padrão.Cada linha na tabela
match_resultsrepresenta 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.
-
Obtenha o bloco CIDR da VPC do workspace do Realtime Compute for Apache Flink.
Faça login no console do Realtime Compute for Apache Flink.
Na coluna Actions do workspace desejado, escolha .
Na caixa de diálogo Workspace details, visualize o CIDR block do vSwitch do Flink totalmente gerenciado.
-
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
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.
-
Adicione
flink-cepcomo 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> -
Desenvolva o código do job.
-
Construa um Kafka Source.
Para detalhes sobre como escrever o código, consulte Kafka DataStream Connector.
-
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> inputO fluxo de eventos de entrada.
PatternProcessorDiscovererFactory<T> discovererFactoryUma fábrica que cria um
PatternProcessorDiscoverer. Este descobridor busca as regras mais recentes e constrói as instâncias correspondentes dePatternProcessor.TimeBehaviour timeBehaviourDescreve 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> outTypeInfoDescreve 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. UmPatternProcessorcontém umPatternespecífico que descreve como corresponder eventos e umaPatternProcessFunctionque descreve como lidar com uma correspondência, como enviar um alerta. Ele também inclui umide umaversionpara identificar oPatternProcessor. Para mais contexto, consulte a proposta.O
patternProcessorDiscovererFactorycria um descobridor para buscar oPatternProcessormais recente. O código de exemplo inclui uma classe abstrata que mostra como consultar periodicamente um armazenamento externo em busca de novas instâncias dePatternProcessor.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
jdbcUrlA URL JDBC do banco de dados.
jdbcDriverO nome da classe do driver do banco de dados.
tableNameO nome da tabela do banco de dados.
initialPatternProcessorsO
PatternProcessorpadrão a ser usado quando a tabela de regras no banco de dados estiver vazia.intervalMillisO 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"); } }NotaPara fins de demonstração, o código de exemplo agrupa o fluxo de dados de entrada por
ideproduct_idantes de conectá-lo aoCEP.dynamicPatterns(). Isso significa que apenas eventos com o mesmoideproduct_idsão considerados para correspondência de regras. Eventos com chaves diferentes não serão comparados entre si. -
-
-
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.
NotaComo 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 falseA 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.NotaO 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.
-
-
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 aoaviatorsã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-firstPara mais informações sobre como configurar parâmetros de tempo de execução, consulte Configuração de parâmetros de tempo de execução.
-
Na página , 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.
Faça login no console do ApsaraDB RDS for MySQL.
-
Insira a regra de atualização dinâmica.
Concatene a string JSON com o
id, aversione 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
patternno 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 campopatternna 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 umaactionigual a 0, a ação do próximo evento não é 1.NotaNo 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.
-
-
Envie mensagens para o tópico
demo_topicusando 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,1662022780000Selecione o método de envio Console, insira
1no 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.
-
-
Visualize a regra mais recente nos logs do JobManager e a correspondência nos logs do TaskManager.
-
Nos logs do JobManager, pesquise por
JDBCPeriodicPatternProcessorDiscovererpara 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 updatedconfirma 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.
-
-
Consulte a tabela
match_resultsexecutandoSELECT * FROMpara ver os resultados que correspondem à regra.A consulta retorna um registro com os campos
rule_id,rule_version,user_id,user_nameeproduction_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.
-
Defina o parâmetro
usingEventTimecomotrue.Na página , localize a implantação desejada e clique em Cancel na coluna Actions.
Em , clique em Edit, defina o parâmetro
usingEventTimecomotruee clique em Save.Clique em Start para iniciar o job novamente.
-
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'); -
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 -
Consulte a tabela
match_resultsexecutandoSELECT * FROMpara ver os resultados que correspondem à regra.A consulta retorna dois registros, com colunas para
rule_id,rule_version,user_id,user_nameeproduction_id. O primeiro registro é para Ken (rule_version=1) e o segundo é para Tom (rule_version=2), ambos associados aoproduction_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.