大数据(9g)FlinkCEP
admin
2024-03-15 19:49:05
0

文章目录

  • 概述
  • 示例代码
    • 环境和依赖
    • Java代码
      • 上面代码可改成下面

概述

  • CEP
    Complex Event Processing:复合事件处理
    通过分析事件间的关系,从事件流中查询出符合要求的事件序列
  • 例如【切菜=>洗菜=>炒菜】3个事件按时间序串联,是正常的事件流
    当发现【切菜=>炒菜】忽略洗菜的事件流,可认为是异常事件

示例代码

环境和依赖

WIN10+JDK1.8+IDEA2021+Maven3.6.3
CEP额外依赖为flink-cep

881.14.62.121.18.24


org.apache.flinkflink-java${flink.version}org.apache.flinkflink-streaming-java_${scala.binary.version}${flink.version}org.apache.flinkflink-clients_${scala.binary.version}${flink.version}org.apache.flinkflink-runtime-web_${scala.binary.version}${flink.version}org.apache.flinkflink-cep_${scala.binary.version}${flink.version}

Java代码

监测 严格近邻的连续三次a的事件流

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;public class CepPractice {public static void main(String[] args) throws Exception {//创建环境StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment().setParallelism(1);//添加数据源,确定水位线策略SingleOutputStreamOperator d = env.fromElements("c", "a", "a", "a", "a", "b", "a", "a").assignTimestampsAndWatermarks(WatermarkStrategy.forMonotonousTimestamps().withTimestampAssigner((element, recordTimestamp) -> 1L));//定义模式Pattern p = Pattern.begin("first").where(new SimpleCondition() {@Overridepublic boolean filter(String value) {return value.equals("a");}}).next("second").where(new SimpleCondition() {@Overridepublic boolean filter(String value) {return value.equals("a");}}).next("third").where(new SimpleCondition() {@Overridepublic boolean filter(String value) {return value.equals("a");}});//在流上匹配模型PatternStream patternStream = CEP.pattern(d, p);//使用select方法将匹配到的事件流取出patternStream.select((PatternSelectFunction) map -> {//Map的key是事件名称(上面的first、second和third)//Map的key对应的value是列表,储存匹配到的事件String first = map.get("first").toString();String second = map.get("second").toString();String third = map.get("third").toString();return first + "->" + second + "->" + third;}).print();//执行env.execute();}
}

打印结果

[a]->[a]->[a]
[a]->[a]->[a]

上面代码可改成下面

留意.times(3).consecutive()

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.List;public class CepPractice2 {public static void main(String[] args) throws Exception {//创建环境StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment().setParallelism(1);//添加数据源,确定水位线策略SingleOutputStreamOperator> d = env.fromElements(Tuple2.of("a", 1000L), Tuple2.of("a", 2000L), Tuple2.of("a", 3000L),Tuple2.of("a", 4000L), Tuple2.of("b", 5000L), Tuple2.of("a", 6000L)).assignTimestampsAndWatermarks(WatermarkStrategy.>forMonotonousTimestamps().withTimestampAssigner((element, recordTimestamp) -> element.f1));//定义模式Pattern, Tuple2> p = Pattern.>begin("=a").where(new SimpleCondition>() {@Overridepublic boolean filter(Tuple2 value) {return value.f0.equals("a");}}).times(3).consecutive(); //严格连续//在流上匹配模型PatternStream> patternStream = CEP.pattern(d, p);//使用select方法将匹配到的事件流取出patternStream.select((PatternSelectFunction, String>) map -> {//Map的key是事件名称(上面的first、second和third)//Map的key对应的value是列表,储存匹配到的事件List> ls = map.get("=a");String first = ls.get(0).f0;String second = ls.get(1).f0;String third = ls.get(2).f0;return String.join("=>", first, second, third);}).print();//执行env.execute();}
}

相关内容

热门资讯

linux入门---制作进度条 了解缓冲区 我们首先来看看下面的操作: 我们首先创建了一个文件并在这个文件里面添加了...
C++ 机房预约系统(六):学... 8、 学生模块 8.1 学生子菜单、登录和注销 实现步骤: 在Student.cpp的...
A.机器学习入门算法(三):基... 机器学习算法(三):K近邻(k-nearest neigh...
数字温湿度传感器DHT11模块... 模块实例https://blog.csdn.net/qq_38393591/article/deta...
有限元三角形单元的等效节点力 文章目录前言一、重新复习一下有限元三角形单元的理论1、三角形单元的形函数(Nÿ...
Redis 所有支持的数据结构... Redis 是一种开源的基于键值对存储的 NoSQL 数据库,支持多种数据结构。以下是...
win下pytorch安装—c... 安装目录一、cuda安装1.1、cuda版本选择1.2、下载安装二、cudnn安装三、pytorch...
MySQL基础-多表查询 文章目录MySQL基础-多表查询一、案例及引入1、基础概念2、笛卡尔积的理解二、多表查询的分类1、等...
keil调试专题篇 调试的前提是需要连接调试器比如STLINK。 然后点击菜单或者快捷图标均可进入调试模式。 如果前面...
MATLAB | 全网最详细网... 一篇超超超长,超超超全面网络图绘制教程,本篇基本能讲清楚所有绘制要点&#...
IHome主页 - 让你的浏览... 随着互联网的发展,人们越来越离不开浏览器了。每天上班、学习、娱乐,浏览器...
TCP 协议 一、TCP 协议概念 TCP即传输控制协议(Transmission Control ...
营业执照的经营范围有哪些 营业执照的经营范围有哪些 经营范围是指企业可以从事的生产经营与服务项目,是进行公司注册...
C++ 可变体(variant... 一、可变体(variant) 基础用法 Union的问题: 无法知道当前使用的类型是什...
血压计语音芯片,电子医疗设备声... 语音电子血压计是带有语音提示功能的电子血压计,测量前至测量结果全程语音播报࿰...
MySQL OCP888题解0... 文章目录1、原题1.1、英文原题1.2、答案2、题目解析2.1、题干解析2.2、选项解析3、知识点3...
【2023-Pytorch-检... (肆十二想说的一些话)Yolo这个系列我们已经更新了大概一年的时间,现在基本的流程也走走通了,包含数...
实战项目:保险行业用户分类 这里写目录标题1、项目介绍1.1 行业背景1.2 数据介绍2、代码实现导入数据探索数据处理列标签名异...
记录--我在前端干工地(thr... 这里给大家分享我在网上总结出来的一些知识,希望对大家有所帮助 前段时间接触了Th...
43 openEuler搭建A... 文章目录43 openEuler搭建Apache服务器-配置文件说明和管理模块43.1 配置文件说明...