ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

基于Flink的实时风控系统实战:规则引擎、状态管理与数据集成全解析

基于Flink的实时风控系统实战:规则引擎、状态管理与数据集成全解析 1. 项目背景与整体设计思路先交代一下我做这个项目的背景。当时团队接到的业务诉求很直白现有交易系统里有一批风控规则跑在离线数仓上T1出结果很多欺诈行为要等第二天才能被发现黑产早就把羊毛薅完了。业务方明确要求至少把一部分高风险场景的识别时效压缩到分钟级甚至秒级。在一轮技术选型之后我们最终把所有实时计算相关的工作都压在了Flink上基于这套引擎搭了一套比较完整的实时风控系统。做这个系统之前我先把风控场景的几个关键特点捋了一遍。第一数据源多且杂既有业务库的订单变更也有前端埋点的行为日志还有第三方的黑名单接口返回这些数据格式、时效性、写入频率完全不一样。第二规则迭代极快风控同学今天可能刚上线一个规则明天就要调整阈值后面还要临时加一个联合规则所以规则层必须和计算层解耦不能每次改规则都重新发布整个作业。第三对延迟非常敏感但又不能为了低延迟牺牲太多准确性需要在实时性和精确性之间做一个可配置的平衡。第四出了问题要能排查风控结果是要回溯的你不能跟审计说“数据丢了查不了”。基于这些约束我选了Flink作为整个系统的底座。理由其实不难理解Flink天然的流处理能力支持毫秒级延迟精确一次语义能保证数据不重不丢加上它强大的状态管理、窗口机制和丰富的连接器生态几乎就是为风控这类场景量身定做的。而且Flink可以同时处理流批两种模式后续做数据回补、模型训练样本抽取都会方便很多。整个系统的架构我最终分成了四层接入层、计算层、存储层、决策层。接入层负责把Kafka里的各种主题数据解析成统一的事件模型计算层跑规则引擎、特征计算和关系图谱分析存储层用Doris加Redis加HBase的组合分别承担明细查询、实时特征缓存和图存储的职责决策层则根据规则命中的结果给出放行、人工审核、拒绝等动作。这里我想重点说说为什么存储层要拆成三个组件而不是一个数据库搞定。Doris适合做大规模明细查询和分析但是单笔查询延迟在毫秒到几十毫秒之间Redis适合做key-value类的实时特征查询抗住高并发没问题但如果要按复杂条件过滤就力不从心了HBase用来存图关系数据比如设备与用户、用户与订单之间的多跳关系。三者各司其职才能同时满足在线决策的高吞吐低延迟和事后分析的灵活性。2. 规则引擎与计算层实现细节2.1 规则引擎选型Flink CEP还是自研表达式引擎这是整个项目里最有争议的一个技术选型。最开始不少同事倾向于直接用Flink CEP因为Flink CEP写起来灵活能做时间窗口内的复杂事件序列匹配比如“5分钟内同一设备登录超过3个账号”这种场景CEP天然支持。但实际做下来我发现了不少问题。首先是规则上线效率风控同学改一条规则我们得改代码、打包、发布、重启作业一次流程快则半小时慢则半天其次是规则总条数多了之后单个CEP作业的复杂度急剧上升调试非常痛苦第三个是CEP的状态很难清理窗口一多状态膨胀很快对内存压力很大。最终我们做了一个折中方案用Flink SQL加自研的轻量级表达式引擎来承接大部分规则只有极少数真正需要复杂序列匹配的场景才用CEP。表达式引擎把风控同学配置的规则编译成可执行的逻辑比如“交易金额大于1000且商户不在白名单中”就转成一段Groovy脚本在Flink的map算子里面去跑。这样改规则只需要改配置下发不需要动Flink作业本身。这里有个经验供各位参考Flink SQL适合做基于时间窗口的聚合类特征比如“近5分钟同一IP的支付次数”这种场景用SQL写非常简洁TUMBLE窗口或者HOP窗口一开就行了。而表达式引擎适合做单条事件的规则命中判断逻辑简单、变更频繁。两者配合才能既保证性能又能快速响应业务变化。2.2 Watermark在风控场景中的设置策略Flink SQL里Watermark的配置直接决定了数据延迟和准确性之间的取舍。风控数据里最典型的问题是用户先做了支付动作然后才上报了登录日志但是支付事件先到了Kafka登录事件后到这会导致Join对不上。如果不做任何处理直接按事件时间关联那么支付事件永远等不到它的前置登录事件。解决这个问题就得靠Watermark延后触发窗口计算用一个forBoundedOutOfOrderness去容忍一定程度的乱序。我给的配置策略是根据数据源的重要程度区分处理。登录、注册这一类高价值事件的延迟容忍度给到10到15秒交易类事件给到5秒左右。这个值不能拍脑袋需要结合业务方确认数据链路的最大延迟周期。一开始我把所有事件的容忍度都设成30秒结果导致风控判定结果整体晚了半分钟被业务方吐槽反应太慢。后来压到5秒之后误杀率上升了一些因为确实有极少量的网络延迟导致事件被丢在了窗口外面。最后我们做成配置化每个事件类型独立设置容忍时间并且在规则引擎里加了一个补偿逻辑对于迟到的数据如果它在规则命中结果已经产出后才到达就发一条补偿事件给下游做二次处理。还有一个细节容易被忽视Watermark不是定义好就万事大吉的它必须在源表中声明并且所有用到事件时间的SQL操作都要显式指定时间字段。很多新人在Flink SQL里配Watermark时把时间字段和Watermark字段混在一起导致下游JOIN根本不会按照预期工作而且这个问题还不会报错排查起来非常折磨人。2.3 状态管理与Checkpoint配置实战Flink的State是实时风控的核心资产。比如你要计算“这个用户近24小时内的累计交易金额”这个累计值就必须存在状态里。传统做法是用Redis来存但Flink本身的状态机制更值得优先考虑。把它理解为每个算子内部的持久化local cache由Flink负责备份恢复比你自己在外面维护Redis再处理缓存失效和一致性要省心得多。我的核心状态配置大概是这样启用Checkpoint间隔设60秒使用增量Checkpoint模式状态后端用RocksDB。选择RocksDB的原因很简单风控系统的状态量级不是KB级别而是GB级别往上走的内存根本扛不住必须落盘。有人问用RocksDB会不会拖慢速度确实会但可以通过调整RocksDB的block cache大小、开启状态压缩来优化。在压测环境里我把block cache配到256MB状态读写性能比默认配置提升了将近40%。每个状态都必须设置TTL。风控特征是有时效性的比如“短时间内登录失败次数”这个状态如果用户上次登录失败发生在三天前那这个状态对当前判断毫无意义白白占着存储资源。我给这类状态统一设了24小时的TTL只有少数真正需要跨天结算的特征才设到7天。TTL实现的时候要注意Flink的过期数据不会立刻被清理而是等到访问或者Compaction的时候才会触发清理所以如果你发现状态文件比预期大别慌这属于正常现象。2.4 双流Join实战设备指纹流与登录行为流双流Join是我在做这个系统时踩坑最多的部分没有之一。简单说一下业务场景一条登录事件到达之后我们要去关联这个设备ID在最近5分钟内的历史登录行为以判断这个设备是否在被多个账号轮番使用。这里就涉及到两条流的Join——一条是实时到达的登录事件流另一条是同设备的历史登录记录流。Flink SQL对双流Join的支持默认是内连接只有左右两边都满足条件才会输出。但是在风控场景里我更常用的是维表关联加上窗口内的增量匹配。更具体的做法是把历史登录记录缓存在状态里然后用interval join把两条流对齐到同一个5分钟窗口内进行关联。interval join的好处是可以精确控制关联的时间范围避免无限期等待另一条流的数据。语法上也简单一条SQL就能搞定SELECT a.device_id, a.user_id, b.login_time FROM login_event AS a JOIN login_history AS b ON a.device_id b.device_id AND a.event_time BETWEEN b.event_time AND b.event_time INTERVAL 5 MINUTE这个场景里最关键的是空窗期的处理。设备第一次出现时Join结果为空此时策略应该是记录这条登录事件并初始化状态而不是直接丢弃。我们利用左连接加COALESCE来兜底保证设备首次登录也能进入规则判断流程。3. 数据同步与外部存储集成的坑与方案3.1 Flink CDC同步业务库数据的注意事项实时风控需要第一时间感知业务库的变化比如订单状态从“待支付”变成“已支付”这条变更如果等离线同步延迟就太大了。Flink CDC组件可以直接监听MySQL、PostgreSQL的binlog或者WAL日志把这个变更流实时推到Kafka再由Flink消费处理。这个东西我用了很久稳定性整体不错但有几个坑必须提醒大家。第一个坑是数据库压力。CDCE把整个库的变更都读出来如果业务库本身压力大binlog增长很快CDC任务的Lag很容易飙升到分钟级。我的解决办法是尽量只同步必要的表加上过滤条件同时把Flink CDC的并行度控制在合理范围内不要一上来就开十几个线程去读同一个库的binlog那是对数据库的变相攻击。第二个坑是类型映射。CDC读出来的数据类型和Flink内部类型有对应关系但是MySQL的JSON类型、Decimal类型在Flink SQL里的处理表现跟你想的往往不完全一致。尤其是Decimal类型默认的精度和scale可能被截断导致下游计算精度丢失。一条交易金额的精度如果出了问题风控规则命中就会出现严重错误这个事非常危险。所以用CDC同步进来之后我强烈建议做一层清洗转换把所有关键字段的精度、类型显式指定好不要依赖默认行为。3.2 Apache Doris类型映射问题排查datev2与dateday这里要分享一个真实的线上事故。我们在用Flink写数据到Doris时作业跑了一段时间突然报错错误信息是type is datev2, but arrow type is dateday. at org.apache.doris.flink.这个报错的意思是Flink这边算出来的字段类型是DATEV2但是Doris连接器期望的Arrow类型是DateDay两者没有匹配上。这个问题的根源在于Flink的Doris连接器在转换数据类型时对DATE类型的映射规则和Doris表的实际类型不一致。当上游源表字段是DATE类型而写入目标是Doris的DATEV2列时连接器生成Arrow RecordBatch时用的类型映射跟Doris端期待的Arrow类型对不上直接就挂了。解决办法有几种我当时用的是第一种。第一在Flink SQL里对源数据做CAST显式把字段类型转成Doris兼容的类型第二升级Doris连接器的版本较新版本对DATE类型映射做了优化第三改Doris表结构把DATEV2类型改成DATE类型让两边对齐。我后来在生产环境同时用了第一和第三种方法问题彻底解决。这里分享的思路是遇到连接器类型报错第一时间先去看连接器版本和存储端版本之间的兼容矩阵别急着改业务逻辑。很多时候是版本之间的类型系统没对齐造成的调整一个不起眼的配置就能解决。3.3 Flink JDBC连接器异常排查实录JDBC连接器在使用过程中遇到的坑也比较多而且很多异常信息写得非常隐晦。我整理一个高频异常清单方便大家直接对号入座。常见的Connection is not available, request timed out一般都指向连接池太小或者数据库侧慢查询阻塞了连接释放。解决办法是调大连接池参数和超时时间另外检查一下目标库有没有长时间占锁的会话把连接资源都耗尽了。我第一次遇到这个报错时以为是Flink并发太高把连接池从10调到100结果数据库直接被压垮了。后来用连接池监控才发现其实是一条慢SQL把数据库连接全占住了根源上得优化查询条件。Communications link failure大概率是网络层抖动或者数据库主动断开了空闲连接。Flink这边需要开启连接自动重连机制同时设置合理的testConnectionOnCheckin参数。你也可以在连接串里加上autoReconnecttrue但这只能兜底网络闪断救不了其他问题。还有一种非常容易误导人的异常报错信息是No suitable driver found。明明本地跑得通提交到集群就报这个。原因是Flink集群的lib目录里没有打包相关的JDBC驱动或者打包时驱动的scope设置不对导致驱动类没有打进去。如果你用的是SQL Client还需要显式地去加载驱动包。解决思路就是检查驱动JAR是否在各节点上实际存在光在pom.xml里加了依赖是不够的。3.4 Flink一定要配HDFS吗存储层设计经验很多新手在学习Flink的时候都有一个疑问Flink是不是一定得装HDFS我的答案是分情况。Flink的核心计算引擎本身完全不依赖HDFS它跑在本地文件系统上也能正常运行。HDFS主要扮演两个角色要么是Checkpoint和Savepoint的存储介质要么是历史数据读写的外部存储。如果你只是做本地测试或者简单的流处理任务用本地文件系统就能把Checkpoint存下来。但生产环境我强烈建议把Checkpoint放到分布式存储上HDFS、S3、OSS都可以否则你的Flink集群一旦发生故障迁移新节点无法从旧节点本地磁盘读取Checkpoint等于你的状态恢复能力直接归零。我们团队因为基础设施里面没有现成的HDFS集群所以最初用了一段时间的本地存储Checkpoint后来在一次节点宕机事故中吃了个大亏整个作业从最近一次Checkpoint恢复结果所有节点本地文件都拿不到了只能从零开始重新消费Kafka数据。最终我们接入了对象存储服务作为Checkpoint存储吞吐和稳定性都不错。所以如果问我“Flink是不是一定要配HDFS”我会说本地测试不一定生产环境要有满足Flink状态持久化能力的外部存储具体用什么得结合你的已有存储底座来选。4. 运维经验与常见问题排查技巧实录4.1 反压问题排查与资源规划心得线上跑了一阵子之后Flink UI上的反压告警开始频繁出现这个反压是Flink最经典的问题之一。用通俗的话讲反压就是下游处理不过来把上游的“路”堵住了。Kafka消费速度变慢消息堆积越来越严重规则判断的时效性也就无从谈起。我遇到的最常见反压源头是规则引擎里跑了一个非常耗时的外部服务调用。当时风控同学希望每个交易事件都实时查一下第三方黑名单结果这个第三方接口的P99响应时间从100毫秒恶化到2秒直接拖垮了整个作业。后续优化方案是给这个调用加上异步IO让Flink不用一条一条阻塞等待外部响应同时给外部调用加上超时熔断机制。Flink Async I/O的吞吐量提升很明显使用之后作业的吞吐比原来提升了将近三倍。资源规划方面我的经验是宁可先小后大不要一开始就按峰值分配。Flink的TaskManager数量可以动态调整但是需要重启作业才能生效。如果你实在估不准资源量可以先给一个相对保守的配置跑几天观察每个算子的繁忙程度和Kafka Lag走势再决定是否扩容。对了还有一个小提示TaskManager的容器内存不要只按堆内内存来算RocksDB堆外内存、网络缓冲区和JVM元空间都要算进去否则跑着跑着容器就被OOM Killed了。4.2 Flink CDC的版本兼容与Docker部署实测项目上线之后我为了给业务方快速搭建一套测试环境决定用Docker Compose把Flink和CDC相关组件一起拉起来。当时Flink主版本已经发布了比较新的版本CDC也发布了大版本两者合在一起的时候出现了不少兼容性问题。比如Flink 2.2.1配Flink CDC 3.5.0之后DDL同步一直报错最后查了官方文档确认是CDC新版本要求Flink的某些依赖包必须单独引入常规的Flink发行版没有预置。Docker部署方式确实能大幅降低环境搭建成本。用一条命令就能把Flink JobManager、TaskManager、MySQL以及CDC服务全部起起来。但要注意容器的资源隔离问题Docker默认不会限制容器使用的CPU和内存如果你只设置了内存上限而没有设置CPU限制多个容器会争抢CPU资源导致Flink作业性能出现剧烈波动。我当时的做法是在docker-compose.yml里给TaskManager配置了CPU配额实测定下来稳定很多。给一个建议Docker环境适合做功能验证和开发调试生产环境尽量还是用原生的集群部署方式资源隔离性和排障便利性都要好得多。docker部署还有个坑是网络模式默认bridge模式下容器之间用服务名互联没问题但一旦任务里配置了外部Kafka地址容器内部不能用localhost去连宿主机上的Kafka得把地址改成宿主机局域网IP才能连通。4.3 数据血缘追踪风控审计的必要基础风控系统的审计要求决定了我们必须要能追踪到每一条规则命中结果的完整链路。比如某笔交易被拒绝业务方来问为什么你要能回答出来是哪个规则命中的依赖了哪些特征这些特征又是从哪些原始数据计算出来的。这个诉求在实时计算场景里做起来比离线要难一些Flink作业是一个长期运行的流式管道中间状态不断变化不像离线任务有明确的调度时间和输入输出关系。我们做了两件事来解决这个问题。第一在规则的输出结果表里增加若干信息列包括命中的规则ID、规则版本号、特征计算时间、参与计算的特征明细的JSON快照。这样每次命中都被完整记录下来事后查起来非常清晰。第二在Flink作业的每个关键算子接入统一的日志切面把算子的输入输出都打上traceId再把这个traceId贯穿到下游Doris的明细表里。通过traceId我可以从一笔交易的最终决策结果一路定位到它的每一层数据来源。数据血缘这件事一开始做觉得繁琐但做过一次之后就明白了它的价值。尤其是规则发生误杀需要回溯的时候你就知道数据链路完整记录到底有多重要了。4.4 常见问题排查速查表我在整个项目落地过程中前前后后处理过不少问题这里整理一个速查表直接按症状定位能帮你省下不少排查时间。现象可能原因解决方案作业启动后一直处于RUNNING但不消费数据并行度过低或Kafka分区数据倾斜调大并行度检查数据Key分布必要时加随机前缀数据有延迟Kafka Lag持续增长反压严重或外部服务调用慢定位反压源头算子使用Async I/O优化外部调用Checkpoint失败但作业不失败状态过大或外部存储吞吐不足调大Checkpoint超时时间优化RocksDB配置检查存储IO结果数据写入Doris报类型错误字段类型映射不一致CAST显式转类型升级连接器核对Doris SchemaFlink SQL中Join数据输出为空时间字段未使用事件时间或晚到数据被丢弃检查Watermark策略调整容忍延迟时间排查时间字段选择规则修改后不生效规则缓存未刷新或版本未更新检查分布式缓存刷新机制确认规则版本号有没有下发成功这些问题的排查过程里最值得关注的是定位思路。多数问题不是一下子就能看出原因的我的习惯是先看Flink UI上的指标曲线再看日志最后才改代码。很多朋友一遇到问题就喜欢马上去翻源码或者改代码结果浪费了很多不必要的时间。5. 风控系统的总结与进阶方向前前后后花了三个多月这套基于Flink的实时风控系统才算是稳定跑在生产环境上。现在每天处理千万级的事件流支持几百条规则的在线运行总体的端到端延迟控制在10秒以内。从离线T1到实时秒级整个风控时效性提升的效果非常明显业务方也确实感受到了变化。踩过这么多坑之后我最大的感受是Flink本身的技术栈其实没多深的水真正的难点全在工程化细节里。规则引擎的选型会影响你后续每一次的规则迭代效率Watermark的参数影响着你系统的准确率和召回率状态的管理影响着你的资源成本和稳定性数据链路的血缘记录影响着审计和排障能不能做下去。这些看起来都不起眼但是放到一起就是能不能支撑住风控业务长期演进的差别。如果这个系统还要继续演进我个人认为两个方向值得投入。第一是引入实时特征平台把特征计算从规则引擎中进一步解耦做成独立的特征服务这样不同规则之间可以复用特征不用每条规则重复计算。第二是把机器学习模型引入实时链路比如在Flink作业里加载一个训练好的XGBoost或者深度学习模型直接把模型打分结果作为规则的一部分。这两块现在业内都有不少成熟实践基于目前的这套底层设施去扩展我觉得不算难。最后再分享一个小技巧写Flink SQL的时候如果遇到结果和预期不符先在SQL Client里把SQL单独跑一遍不要直接丢到作业里去排查。SQL Client会给你非常直接的执行计划和错误信息比在集群上盯着日志看要高效得多。这个习惯很多老手都在用新手往往容易忽略但其实非常实用。
返回列表