1. 项目背景与整体设计思路
先交代一下我做这个项目的背景。当时团队接到的业务诉求很直白:现有交易系统里有一批风控规则跑在离线数仓上,T+1出结果,很多欺诈行为要等第二天才能被发现,黑产早就把羊毛薅完了。业务方明确要求,至少把一部分高风险场景的识别时效压缩到分钟级甚至秒级。在一轮技术选型之后,我们最终把所有实时计算相关的工作都压在了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参数。你也可以在连接串里加上autoReconnect=true,但这只能兜底网络闪断,救不了其他问题。
还有一种非常容易误导人的异常,报错信息是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 Schema |
| Flink SQL中Join数据输出为空 | 时间字段未使用事件时间或晚到数据被丢弃 | 检查Watermark策略,调整容忍延迟时间,排查时间字段选择 |
| 规则修改后不生效 | 规则缓存未刷新或版本未更新 | 检查分布式缓存刷新机制,确认规则版本号有没有下发成功 |
这些问题的排查过程里,最值得关注的是定位思路。多数问题不是一下子就能看出原因的,我的习惯是先看Flink UI上的指标曲线,再看日志,最后才改代码。很多朋友一遇到问题就喜欢马上去翻源码或者改代码,结果浪费了很多不必要的时间。
5. 风控系统的总结与进阶方向
前前后后花了三个多月,这套基于Flink的实时风控系统才算是稳定跑在生产环境上。现在每天处理千万级的事件流,支持几百条规则的在线运行,总体的端到端延迟控制在10秒以内。从离线T+1到实时秒级,整个风控时效性提升的效果非常明显,业务方也确实感受到了变化。
踩过这么多坑之后,我最大的感受是:Flink本身的技术栈其实没多深的水,真正的难点全在工程化细节里。规则引擎的选型会影响你后续每一次的规则迭代效率;Watermark的参数影响着你系统的准确率和召回率;状态的管理影响着你的资源成本和稳定性;数据链路的血缘记录影响着审计和排障能不能做下去。这些看起来都不起眼,但是放到一起就是能不能支撑住风控业务长期演进的差别。
如果这个系统还要继续演进,我个人认为两个方向值得投入。第一是引入实时特征平台,把特征计算从规则引擎中进一步解耦,做成独立的特征服务,这样不同规则之间可以复用特征,不用每条规则重复计算。第二是把机器学习模型引入实时链路,比如在Flink作业里加载一个训练好的XGBoost或者深度学习模型,直接把模型打分结果作为规则的一部分。这两块现在业内都有不少成熟实践,基于目前的这套底层设施去扩展,我觉得不算难。
最后再分享一个小技巧:写Flink SQL的时候,如果遇到结果和预期不符,先在SQL Client里把SQL单独跑一遍,不要直接丢到作业里去排查。SQL Client会给你非常直接的执行计划和错误信息,比在集群上盯着日志看要高效得多。这个习惯很多老手都在用,新手往往容易忽略,但其实非常实用。