Flink实时告警系统设计与实现
Flink实时告警系统设计与开发
背景
实时监控系统需要满足对多种来源的数据进行告警,为提升系统的可扩展性和灵活性,采用动态规则配置来实现多种数据源、多种告警规则的实时告警。需要实时监测和发现车端云端的信号、埋点数据是否有异常,车辆运行状况异常。
1. 数据来源

2. 系统架构设计
1. 系统分层架构设计

本着高内聚低耦合的原则,实时告警系统采用分层设计的思想对整体的功能模块进行组合,其中:
- Flink DataStream 层:数据流在Flink内部的整体流向DAG图,如
addSource、connect、process、addSink。 - Flink Function 层:对function的具体实现,如
AlertManagerSinkFunction、CustomMysqlSourceFunction、RuleMatchBroadCastProcessFunction等。 - Service 层:业务的处理过程,如负责向AlertManager传输数据的
AlertManagerService、负责规则同步、更新、维护、转化、匹配的RulesService。
2. 业务模块设计


说明:业务上,需要告警的数据源目前有4种数据来源,分别是远端日志、云端微服务日志、车机端埋点、Sentry异常崩溃。其中Sentry中的数据需要通过告警规则的筛选后发送到Kafka中用于实时监控。设计上首先通过Driver中的class路由到通用JSON告警模块或者Sentry异常崩溃业务处理模块,其次通过app.type选择Kafka中的数据源。
3. Flink DataStream 处理流程图

说明:DataStream处理流程图展示的是数据从Kafka消费后在Flink Function中的流向关系。Driver负责Flink程序的启动,通过class筛选路由到通用JSON告警或者Sentry异常崩溃模块,其中内部的逻辑比较相似:
- 首先Mysql中的配置通过自定义数据源模块会被解析成配置流。
- 其次Kafka topic会被解析成数据流,通过广播连接,配置流会被广播到每个数据流的TaskManager。
- 通过规则匹配模块对数据流和规则流进行匹配。
- 匹配到数据筛选出非Sentry中的数据分别发送到AlertManager实时告警、MySQL告警统计、Kafka实时监控。
4. 规则引擎使用
Aviator是一个高性能、轻量级的Java语言实现的表达式求值引擎,主要用于各种表达式的动态求值。Aviator是直接将表达式编译成Java字节码,交给JVM去执行。

规则匹配模块核心使用的是Aviator规则引擎表达式进行规则匹配,匹配的内容来源于:
- 数据流的JSON通过
flattenAsMap转成map。 - 规则流中有效的Rule中获取得到的规则表达式。
5. 规则设计
规则存储在MySQL中便于管理和修改,通过Flink CDC可实现动态修改和同步。
end
