ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink实时风控系统落地复盘:特征计算、规则热更新与排障实践

Flink实时风控系统落地复盘:特征计算、规则热更新与排障实践 把Flink接进风控系统之后我最大的一个感悟是实时风控这个事的难点从来不在Flink本身。框架的API、窗口、状态管理熟读文档总能学会真正让团队掉进坑里的是那些藏在实时二字背后的数据对齐、规则热更新、连接器版本匹配和故障恢复问题。这篇文章不是Flink入门教程而是一个基于Flink构建风控系统的复盘记录覆盖从架构设计、特征计算、部署实践到排障调参的完整链路。如果你正在把实时风控从想法推向落地或者已经上线但总觉得哪里不稳这篇文章值得你花点时间看完。1. 实时风控的真实瓶颈为什么是Flink而不是其他框架1.1 风控场景下实时到底意味着什么很多人对实时风控的理解是算得快其实不准确。风控判断一笔交易是否可疑靠的是一系列历史特征这个用户最近5分钟下了多少单、这个IP在过去1小时关联了多少张卡、这个设备指纹最近24小时是否触发过其他规则。这些特征的本质是截至当前时刻业务状态的一种聚合。要计算这种聚合框架必须解决三个问题怎么定义时间、怎么处理乱序数据、怎么在故障时保证状态不丢。这三点恰恰是Flink相对其他框架最核心的优势。Flink把事件时间Event Time作为一等公民用Watermark处理乱序用状态后端保存中间结果用Checkpoint提供故障恢复。比如要算1小时内同一设备登录的不同账号数超过5个在Flink里就是一个带状态的窗口聚合按设备ID分组、开1小时滑动窗口、窗口内统计账号去重数超过阈值直接触发告警。同样的逻辑要是放在老一代流框架里事件时间对齐、窗口清理、故障恢复全都得自己写工程量完全不在一个量级。1.2 为什么不是Spark Streaming或Storm我见过不少团队在风控选型时纠结于Spark Streaming。Spark Streaming的微批模型决定了它的延迟下限通常是秒级这在大多数风控场景里可以接受问题在于它的事件时间处理、状态管理、以及对乱序数据的表达始终没有Flink那么自然。Storm是真正的毫秒级延迟但它太底层了窗口、状态、精确一次都缺乏内置支持开发成本和维护成本都很高。Kafka Streams在简单场景里很香可一旦涉及多作业协同、复杂拓扑、以及和外部存储的一致性写入它的集群能力就比较吃亏。Flink在风控这个场景里胜出的关键不是单纯的速度而是正确的时间语义可靠的状态管理这套组合拳。对于风控这种误判一次可能产生资损、漏判一次可能直接产生坏账的业务来说算得准、故障后能恢复比算得快重要得多。1.3 Flink在风控系统里负责什么不负责什么把Flink的定义域说清楚能避免很多架构争论。我的经验是Flink负责的是实时特征计算、规则判断、以及和外部系统的数据交互它不负责存储海量明细数据也不负责OLAP分析。明细数据归档、事后回溯查询那是数仓和OLAP引擎的事。风控系统的实时主链路里Flink是计算引擎和状态载体而不是万能的数据平台。明确这个边界之后下面聊整体链路时会轻松很多。2. 核心链路拆解接入、特征计算、规则热更新与结果写回2.1 数据接入层Kafka是统一入口别让Flink直接对接几十个数据源风控的数据来源通常包括交易流水、登录日志、设备指纹、埋点行为、用户档案变更等来源系统可能有几十个。如果让Flink作业直接对接每个数据源任何一个上游抖动都会直接影响风控作业的稳定性。我的做法是所有数据先进KafkaFlink只从Kafka消费。Kafka在这里的价值是削峰填谷、解耦、以及最重要的可重放。风控作业做版本升级或者代码回滚时往往需要从某个时间点重新消费数据没有Kafka的消息留存这个操作根本做不了。另外Kafka的Topic分区数也是后续设置Flink并行度的重要依据一个分区对应一个消费线程可以让数据倾斜和并行度管理都变得可控。2.2 实时特征计算的两种模式特征计算是风控规则的核心输入我习惯把它分成两类。第一类是窗口聚合特征比如近5分钟交易金额近1小时登录失败次数当日首次交易距现在的时长。这类特征用Flink SQL的窗口函数就能很优雅地表达开发效率高可读性也强。第二类是跨事件关联特征典型例子是注册手机号与当前交易手机号不一致。这需要把用户注册事件缓存在状态里等交易事件到达后做关联判断。这类逻辑用DataStream API的Keyed State更合适因为涉及具体的状态结构设计和过期策略。这里有个容易犯的错把所有特征都用SQL写。复杂关联用SQL硬写既难调试又难维护。我目前的实践是窗口聚合类特征尽量SQL化跨事件关联类特征用DataStream API封装成可复用的算子两类代码通过统一的特征服务层对外暴露。2.3 规则引擎不要硬编码规则用Broadcast Stream做动态更新风控规则的特点是多变运营同学可能每周都要调整阈值、新增规则。如果把规则写死在作业代码里每次改动都要重新提交作业、恢复状态、验证结果周期太长。我用的是Flink的Broadcast Stream机制规则变更写入配置中心我用的是Nacos一个单独的流实时监听配置变更并广播到所有算子实例业务数据流和规则广播流做connect每条数据就能用最新规则来评估。这里有一个很容易踩的坑规则引用的特征字段和特征计算逻辑之间存在版本兼容问题。规则升级先于特征上线或者特征下线时还有规则引用它都会导致短时间的悬空引用。我的应对办法是在规则结构里带上特征schema版本号特征输出也带版本号连接时做一次版本校验不匹配就丢弃或者走降级分支。这个设计能避免很多线上诡异问题。2.4 结果写回三路输出必须幂等决策结果写到哪里决定了后续整个风控运营体系能不能转起来。我一般把结果分成三路Redis供线上交易链路实时查询拦截结果要求低延迟用决策ID作为key做幂等Doris或ClickHouse供风控运营和数据分析同学做多维分析、案件回溯Kafka回流给下游业务系统做联动处理比如触发二次验证、人工审核。这里的关键词是幂等。Flink作业重启后如果重复写了一条决策结果下游可能因此重复扣款、重复发验证码这是线上事故级别的问题。所以写Redis必须用唯一决策ID做key写Doris这类支持主键的存储要用主键模型做覆盖写。3. 部署落地复盘Flink 2.2.1 与 Flink CDC 3.5.0 的 Docker 实战3.1 版本组合怎么选别默认最新Flink配最新CDC这次项目我用了Flink 2.2.1配合Flink CDC 3.5.0。选这套组合之前我先去确认了CDC官方文档的版本兼容矩阵。这一点特别重要CDC连接器不是Flink内置的需要单独下载jar包放到Flink的lib目录下版本不匹配最常见的报错是NoSuchMethodError或者ClassNotFoundException这种问题排查起来非常痛苦因为它明显不是业务代码的锅但又会让人误以为是代码写错了。顺便说一句Flink的JDBC连接器、Kafka连接器、CDC连接器都是独立于Flink核心的组件它们的版本号对应的Flink版本各有不同。做部署规划时先列一张版本对照表把Flink核心版本、各连接器jar版本、JDBC驱动版本写清楚能省掉后面很多麻烦。3.2 Docker部署的具体步骤和注意点用Docker部署Flink好处是环境一致坏处是资源管理和网络配置要更小心。我这次的部署用docker-compose编排核心服务是JobManager和TaskManager。有几个细节必须注意。第一内存参数要显式设置jobmanager.memory.process.size和taskmanager.memory.process.size都要写清楚否则容器会因为JVM实际使用的内存超过了容器限制而被杀掉尤其在使用RocksDB状态后端时堆外内存的使用量很容易超预期。第二JobManager和TaskManager之间通过RPC通信容器必须放在同一个自定义网络中。第三Checkpoint和Savepoint的目录要挂载宿主机目录不然容器一删所有状态全没了这个坑我见过太多人踩。3.3 Flink一定要HDFS吗这个话题在团队里争论过很多次。我的结论是如果你的部署是多节点的生产集群你需要一个所有TaskManager都能访问的共享存储来放CheckpointHDFS是经典选择但不是唯一选择。如果只是单机测试或者小规模业务Flink的本地文件系统Checkpoint完全可以跑不需要HDFS。但生产环境多台TaskManager各自有本地磁盘Checkpoint写到某一台机器的本地路径其他机器恢复时根本找不到状态文件。这时候需要S3、OSS、Ceph这类共享对象存储或者HDFS。Flink本身不强制依赖HDFS你只需要把对应的Hadoop依赖引入并把Checkpoint路径指过去就行。如果你的公司已经有对象存储优先用对象存储维护成本比自建HDFS低很多。3.4 SQL Client和SQL Gateway的取舍Flink SQL在做风控特征计算时非常好用但它的适用边界要清楚。Flink SQL Client适合本地调试和一次性任务提交不适合给团队里的多个人共享使用。SQL Gateway则是把SQL提交能力做成了服务上层平台可以通过API向Flink集群提交SQL作业算法同学和分析师也能自助提作业不用每次都找平台组开权限。我在风控平台里集成了SQL Gateway但加了审批和资源限制。因为SQL作业和DataStream作业共享同一个集群的资源如果没有配额限制一个写了全表扫描的SQL能把整个风控集群的算力吃掉。SQL Gateway是效率工具但也意味着新的野作业入口治理要跟上。4. 三次真实排障JDBC连接器、Doris类型映射、Watermark不触发4.1 Flink JDBC连接器异常从No suitable driver到连接数打爆热搜词里专门有flink的jdbc连接器异常我猜踩过这个坑的人不少。最常见的报错是长这样的java.sql.SQLException: No suitable driver found for jdbc:mysql://...或者Could not find any factory for identifier jdbc that implements DynamicTableFactory排查链路我整理成四步确认连接器jar在不在Flink lib目录或者作业依赖里。JDBC连接器不是Flink默认自带的必须显式引入flink-connector-jdbc版本要和Flink主版本匹配。确认MySQL驱动jar有没有引入。Flink的JDBC连接器本质上是包装了JDBC驱动但真正实现com.mysql.cj.jdbc.Driver的包在mysql-connector-java里这个驱动包漏掉就会出现No suitable driver。确认网络可达。这个和框架无关但排查顺序经常被忽略。先在TaskManager所在的机器上telnet一下数据库端口很多本地能连、上线就不通的问题根本原因就是安全组或者网络策略没放通。确认连接数没有被耗尽。作业并发度高时每个subtask都会持有JDBC连接数据库连接数很容易被打爆。降低Sink并发或者用支持连接池的配置都能缓解。我在风控项目里JDBC连接器的主要用途是维表关联比如实时查询黑名单库。大批量结果写入不要走JDBC写入性能差还容易拖垮数据库用Doris的Stream Load或者Kafka加下游入库的方式会合理得多。4.2 Doris连接器类型映射DATEV2与DATE的报错处理搜索词里有句报错很典型flink type is datev2, but arrow type is dateday. at org.apache.doris.flink.这个报错发生在Flink通过Doris连接器读写数据时Doris 2.x版本之后日期类型默认是DATEV2而Flink读到的是DATE两边在Arrow序列化协议里的类型对不上。排查这个问题的思路不是去争论谁的bug而是做显式类型映射看报错发生在读还是写。如果是读Doris维表可以在Flink SQL里对日期字段做CAST比如CAST(event_date AS STRING)让类型别在协议层硬碰硬。如果是写Doris在Doris建表时把日期字段类型统一成DATEV2或者干脆用STRING在Flink和Doris之间传入库时由Doris侧转换。检查Flink Doris Connector的版本和Doris服务端版本是否匹配版本差异大时这种底层类型不兼容问题会更频繁。这类问题的本质是分布式系统里的类型系统不一致最后都会落到在两个系统之间做显式转换这个解法上。看到这类报错不用慌把字段类型在两边对齐即可。4.3 Watermark不触发窗口计算三个隐蔽原因还有一个高频问题窗口数据明明已经超过窗口结束时间了但窗口就是不触发输出或者数据一直滞留在窗口里。我在Flink SQL里遇到过一次排查过程很典型。第一步确认事件时间字段类型。如果时间字段是字符串要保证它被正确解析成TIMESTAMP(3)而且Watermark定义要写在事件时间字段上。WATERMARK FOR ts AS ts - INTERVAL 5 SECOND里ts的类型不对Watermark根本推不动。第二步确认数据本身的事件时间有没有在推进。如果上游数据都用了同一个旧时间戳Watermark就会卡住窗口永远不触发。第三步检查乱序容忍时间。容忍时间设得太长窗口计算会延迟很久设得太短很多乱序数据又会被丢弃。风控规则对时效性敏感我一般设3到5秒。第四步检查Kafka Source的空闲分区。Flink消费Kafka时每个分区的Watermark由该分区的数据推进全局Watermark取所有分区的最小值。如果一个分区长时间没有新数据它的Watermark停在初始值全局Watermark就永远不前进。解决办法是配置Source的Idleness Timeout让长时间没有数据的分区不再拖后腿source.assignTimestampsAndWatermarks( WatermarkStrategy.EventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withIdleness(Duration.ofSeconds(10)) );5. 数据一致性CDC链路、端到端精确一次与数据血缘5.1 Flink CDC在风控里的典型用法风控除了流式行为日志还依赖业务库的变更数据比如用户被标记为黑名单、设备被解禁、商户状态变更。这些是低频但高价值的变更用Flink CDC监听MySQL的binlog打成数据流是当前很成熟的方案。Flink CDC 3.x相比2.x的一大变化是把增量快照和整库同步能力大幅增强可以只用YAML配置就完成数据同步任务不用写大量Java代码。比如下面这个配置就能把MySQL的订单变更同步到Kafkasource: type: mysql hostname: localhost port: 3306 username: root password: 123456 tables: order_db.orders server-id: 5400-5404 sink: type: kafka properties.bootstrap.servers: localhost:9092 topic: order_binlog用CDC之前源库的binlog参数必须提前确认[mysqld] server-id123 log_binmysql-bin binlog_formatROW binlog_row_imageFULLbinlog_format必须是ROWbinlog_row_image必须是FULL否则CDC读不到完整的字段前后镜像很多风控需要的变更前值就拿不到。5.2 精确一次的现实选择不要追求所有环节都事务化Flink的Checkpoint机制可以保证算子级别的精确一次但端到端的精确一次依赖每个Sink的支持。Kafka Sink支持精确一次Doris配合Stream Load也支持两阶段提交。但这里我要泼一盆冷水不要在所有环节都追求事务化事务是有成本的会带来更高的延迟和吞吐损失。我的做法是区分场景。高风险交易判断这类核心链路的写入用幂等设计加事务机制保证严格精确一次。辅助决策、特征落库这类场景允许幂等重试下游消费时用唯一ID做去重语义上就能达到近似精确一次。这样能省下大量不必要的性能损耗。5.3 数据血缘在风控场景里的真实价值搜索词里有flink 数据血缘这个功能在风控场景的价值比一般业务系统更大。风控经常要回答一个问题这个用户为什么被限权了答案是哪个规则、哪个特征、哪个数据源触发的。没有血缘关系审计和用户申诉会非常难查。Flink生态里通过解析JobGraph和SQL的字段血缘可以把源表-中间计算-结果表-规则关联起来。我的建议是从SQL作业开始做血缘成本最低DataStream作业的血缘可以先靠代码注释和规范维护。优先把规则维度的血缘打通解决某条规则命中哪些人、由哪些数据产出这个审计刚需比纠结字段级血缘的细粒度更实用。血缘做起来之后再配合规则版本管理风控策略的每一次调整都能追根溯源这在应对合规审计时是实打实的帮助。6. 调参与运维经验并行度、反压、Checkpoint 的取舍6.1 并行度和资源怎么估并行度的设置方法不少我的基准很简单Source并行度以Kafka分区数为参考一个分区对应一个并行度效率最高。特征聚合算子的并行度通常和Key的分布对齐比如按用户ID分组的特征并行度就不宜太低否则单个算子的状态量太大。资源估算的核心是状态大小。估算法是单条状态记录大小乘以状态条目数再考虑复制因子。风控场景里设备到账号关联这种状态很容易做到数百GB。内存不够时RocksDB是兜底方案但要清楚RocksDB会增加CPU开销和GC压力。给Flink容器配内存时显式设置taskmanager.memory.process.size并且给堆外内存留足余量别让容器被系统杀掉。6.2 反压问题的排查逻辑Flink UI上Source反压显示HIGH不代表Source有问题而是下游某个算子处理不过来卡住了整条链路。排查时要顺着拓扑往下找看哪个算子的反压最高。风控作业里最常见的反压源不是CPU密集计算而是外部依赖比如每条数据都同步查一次Redis或者HBase。解决思路有三个查维表改成异步I/O把同步查询变成异步并发请求把变更不频繁的维表做成Broadcast状态避免每条数据都查外部存储给热点维表加本地缓存。我在风控作业里把黑名单维表做成Broadcast状态之后反压直接下降了一个量级这是性价比非常高的优化。6.3 Checkpoint配置的几个经验值Checkpoint间隔和超时时间的设置直接关系到故障恢复的速度。如果业务对恢复时间有要求Checkpoint间隔可以设短一点30秒到1分钟这样故障恢复时只需要回放最近很少的数据不会造成恢复后大量的数据积压和延迟。同时要配好未完成Checkpoint的重试策略避免一个坏掉的Checkpoint把整个作业卡住。Savepoint我要单独提醒一句至少保留最近两个可用的Savepoint。风控规则升级时经常要用Savepoint做状态兼容性验证新代码一旦有问题必须能立刻回滚。只留一个Savepoint遇到状态结构变化时可能根本恢复不了。最后再分享一个我踩了好几次才长记性的经验不管风控规则多复杂先把日志的traceId打通。让一条决策链路从数据进来到结果写出所有环节都能通过同一个traceId串起来。没有这个基础反压定位、数据一致性排查、规则命中归因都会变成猜谜。这项工作的性价比比任何一项框架调优都高。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进