简介:这是一份面向大数据开发初学者与进阶工程师的Flink实战项目资料,基于保险行业真实业务场景,采用Flink+HBase+Kafka+Phoenix技术架构,实现业务系统数据库数据的实时同步与实时统计报表分析,适合希望从理论走向落地、掌握流式计算完整链路的开发者。资源包共73个文件,约584KB,以34个Java源码为核心,辅以11个Python脚本、9个Shell脚本,以及CSV数据样例、properties与XML配置、SQL建表语句、Phoenix初始化文件、Jar依赖和README说明文档,覆盖数据采集、处理、存储到报表输出的完整工程结构。目前已有1166人学习下载,作者另提供Flink答疑服务,便于快速理解项目与入门。通过该资源,读者可参考真实项目的目录组织与代码实现,理解Kafka接入、Flink计算、HBase与Phoenix存储查询的协作方式,并借助脚本与配置快速搭建本地运行环境,积累实时数仓与报表开发的排错经验。
1. 保险保单实时风控:Flink 实战项目到底在算什么
保险行业的数据有个特点:单笔金额不大,但笔数和状态流转极多。一份保单从投保、核保、缴费、批改到理赔,中间会经过十几个状态节点,每个节点都可能触发风控规则。传统做法是 T+1 跑批,第二天早上出报表,等发现异常保单时,钱可能已经赔出去了。Flink 实战项目在保险行业最典型的落地场景,就是把这条链路从「隔夜看」压到「秒级拦」。
这个项目要解决的核心问题有三个:第一,保单事件从业务库出来之后,怎么在秒级内完成规则匹配;第二,规则命中之后,怎么保证不重复告警、不遗漏告警;第三,理赔高峰期流量翻十倍时,作业不能崩。适合谁看?有 Java 或 Scala 基础、写过简单 Flink 作业但没碰过真实业务链路的同学,以及正在做保险、银行、消费金融实时风控的工程师。下面按「数据怎么进来 → 规则怎么算 → 状态怎么存 → 线上怎么稳」的顺序拆开讲。
2. 保单事件接入:从业务库到 Flink 的三种取数方式与选型
保险核心系统多数是 Oracle 或 MySQL,保单状态变更写在业务表里。要把这些变更实时送进 Flink,常见做法有三种:CDC 抓 binlog、业务方发 MQ、定时扫增量表。选哪种,取决于你对延迟和侵入性的容忍度。
2.1 CDC 抓 binlog:延迟最低但最怕 DDL
CDC 方案(Debezium + Kafka 或 Flink CDC)直接读数据库日志,对业务零侵入,延迟可以做到毫秒级。保险场景里,保单表、批单表、理赔表通常都要抓。配置上最关键的是server-id不能和现有从库冲突,snapshot.mode建议用schema_only而不是initial,否则第一次启动会把全表快照拉一遍,几千万行保单能把 Kafka 打满。
# Flink CDC MySQL source 关键配置(放在 SQL DDL 里) CREATE TABLE policy_cdc ( policy_id BIGINT, policy_no STRING, status STRING, premium DECIMAL(18,2), update_time TIMESTAMP(3), PRIMARY KEY (policy_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '10.0.1.20', 'port' = '3306', 'username' = 'flink_reader', 'password' = '******', 'database-name' = 'insurance_core', 'table-name' = 't_policy', 'server-time-zone' = 'Asia/Shanghai', 'scan.incremental.snapshot.enabled' = 'true' );scan.incremental.snapshot.enabled设为 true 之后,快照阶段会加锁读而不是全表锁,对线上业务影响小很多。server-time-zone必须显式指定,否则 binlog 里的时间戳会按 UTC 解析,保单的update_time会差 8 小时,后面做时间窗口全乱。这是血泪经验,第一次上线时规则窗口怎么算都不对,查了两天才发现是时区。
2.2 MQ 接入:业务方主动发消息,可控但依赖改造
如果业务方愿意配合,在保单状态变更的代码里发一条 Kafka 消息,这是最干净的方式。消息体里带上policy_id、event_type、event_time、operator。好处是字段可控、顺序可控,坏处是要推动业务方改代码,保险公司的核心系统往往外包给多家厂商,推动一次改造周期很长。
我一般会建议:新系统直接走 MQ,老系统先用 CDC 兜底,等业务方排期改造后再切。两条链路可以并行跑一段时间做对账,确认 MQ 消息不丢不重之后再下线 CDC。
2.3 定时扫增量表:最土但最稳的兜底方案
有些保险公司的 DBA 不允许开 binlog 给外部消费,业务方也推不动。这时候只能用定时任务扫update_time > last_max_time的增量数据,写进 Kafka 再给 Flink 消费。延迟取决于扫描间隔,一般设 30 秒到 1 分钟。缺点是update_time如果没有索引会全表扫,必须让 DBA 加索引;另外同一秒内多条更新可能因为>和>=的边界问题漏数据,建议用>=加去重。
三种方式没有绝对优劣,保险行业里经常是混合使用:理赔表走 CDC,保单表走 MQ,批单表走定时扫。选型时先问三个问题:DBA 让不让读 binlog、业务方改不改得动、能接受多大延迟。
3. 核保规则实时计算:用 KeyedProcessFunction 做状态化匹配
数据进来之后,核心是规则引擎。保险核保规则少则几十条,多则上千条,包括「同一被保人 30 天内投保超过 3 次」「保费超过年收入 5 倍」「职业类别为拒保职业」等。这些规则需要跨事件的状态,比如「30 天内次数」就必须记住历史。
3.1 为什么不用简单窗口而用 KeyedProcessFunction
窗口能解决「固定时间范围内计数」,但保险规则经常是「自上次理赔之日起 90 天内再次报案」这种滑动起点,窗口的边界对不上。KeyedProcessFunction 可以按policy_id或insured_id做 key,在ValueState里存历史事件列表或计数器,配合Timer做超时清理,灵活性最高。
public class UnderwriteRuleFunction extends KeyedProcessFunction<String, PolicyEvent, Alert> { // 存该被保人最近 30 天的投保时间戳列表 private transient ListState<Long> recentApplyTimes; // 存该被保人累计理赔金额 private transient ValueState<BigDecimal> claimAmount; @Override public void open(Configuration params) { ListStateDescriptor<Long> applyDesc = new ListStateDescriptor<>("apply-times", Long.class); recentApplyTimes = getRuntimeContext().getListState(applyDesc); ValueStateDescriptor<BigDecimal> claimDesc = new ValueStateDescriptor<>("claim-amount", BigDecimal.class); claimAmount = getRuntimeContext().getState(claimDesc); } @Override public void processElement(PolicyEvent event, Context ctx, Collector<Alert> out) throws Exception { String insuredId = event.getInsuredId(); long now = event.getEventTime(); if ("APPLY".equals(event.getEventType())) { // 清理 30 天前的记录 List<Long> valid = new ArrayList<>(); for (Long t : recentApplyTimes.get()) { if (now - t <= 30L * 24 * 3600 * 1000) { valid.add(t); } } valid.add(now); recentApplyTimes.update(valid); if (valid.size() > 3) { out.collect(new Alert(insuredId, "FREQ_APPLY", "30天内投保超过3次,当前=" + valid.size())); } // 注册 30 天后的清理定时器 ctx.timerService().registerEventTimeTimer(now + 30L * 24 * 3600 * 1000); } if ("CLAIM".equals(event.getEventType())) { BigDecimal total = claimAmount.value() == null ? BigDecimal.ZERO : claimAmount.value(); total = total.add(event.getAmount()); claimAmount.update(total); if (total.compareTo(new BigDecimal("500000")) > 0) { out.collect(new Alert(insuredId, "HIGH_CLAIM", "累计理赔超50万,当前=" + total)); } } } @Override public void onTimer(long ts, OnTimerContext ctx, Collector<Alert> out) throws Exception { // 定时器触发时清理过期状态,防止状态无限增长 List<Long> valid = new ArrayList<>(); for (Long t : recentApplyTimes.get()) { if (ts - t <= 30L * 24 * 3600 * 1000) { valid.add(t); } } recentApplyTimes.update(valid); } }ListState存时间戳列表,每次新事件来时先清理过期项再判断。ValueState存累计理赔金额,只增不减,因为理赔金额是历史累计。onTimer里做状态清理,这是防止 RocksDB 状态膨胀的关键——如果不清理,一个被保人跑了三年,recentApplyTimes里会堆几千个时间戳,状态后端迟早撑爆。
3.2 规则热更新:用 BroadcastState 避免重启作业
保险规则经常变,比如监管突然要求「某类职业保额上限下调」。如果规则写死在代码里,每次改规则都要重新打包、上传、重启作业,状态还得从 savepoint 恢复,运维成本高。常见做法是用 BroadcastState:把规则表通过一个低延迟的 Kafka topic 广播到每个并行子任务,规则变更时发一条新规则消息,作业不重启就能生效。
// 广播规则流 DataStream<Rule> ruleStream = env .addSource(new FlinkKafkaConsumer<>("rule_topic", new RuleSchema(), props)) .broadcast(ruleStateDescriptor); // 主流程连接广播流 BroadcastConnectedStream<PolicyEvent, Rule> connected = policyStream.connect(ruleStream); connected.process(new BroadcastProcessFunction<PolicyEvent, Rule, Alert>() { @Override public void processElement(PolicyEvent event, ReadOnlyContext ctx, Collector<Alert> out) throws Exception { ReadOnlyBroadcastState<String, Rule> rules = ctx.getBroadcastState(ruleStateDescriptor); for (Map.Entry<String, Rule> e : rules.immutableEntries()) { if (e.getValue().matches(event)) { out.collect(new Alert(event.getInsuredId(), e.getKey(), e.getValue().getMessage())); } } } @Override public void processBroadcastElement(Rule rule, Context ctx, Collector<Alert> out) throws Exception { ctx.getBroadcastState(ruleStateDescriptor).put(rule.getRuleId(), rule); } });processBroadcastElement里更新规则,processElement里读取规则做匹配。注意 BroadcastState 在每个并行子任务里都有一份完整副本,规则数量不能太大,几百条以内没问题,上万条会占内存。规则匹配的逻辑建议做成可配置的表达式,比如用 Aviator 或 SpEL 解析字符串规则,这样新增规则不用改代码。
3.3 告警去重:用状态加定时器做幂等
同一个被保人可能连续触发同一条规则,比如 30 天内第 4 次投保触发一次,第 5 次又触发一次。如果每次都发告警,风控人员会被淹没。常见做法是在状态里记录「上次告警时间」,同一个insured_id + rule_id在 24 小时内只发一次。
ValueState<Long> lastAlertTime = getRuntimeContext().getState( new ValueStateDescriptor<>("last-alert", Long.class)); Long last = lastAlertTime.value(); if (last == null || now - last > 24L * 3600 * 1000) { out.collect(alert); lastAlertTime.update(now); }这个逻辑简单但有效。注意lastAlertTime也要注册定时器清理,否则长期不活跃的被保人状态会一直留在 RocksDB 里。我一般会设 7 天清理一次,因为保险风控的关注周期通常不超过一周。
4. 状态后端与 Checkpoint:保险场景下的参数怎么调
保险作业的状态量取决于被保人数量。一个中型保险公司活跃被保人几百万,每个被保人的状态几百字节到几 KB,总状态量在几十 GB 到几百 GB 之间。这个量级下,状态后端和 Checkpoint 配置直接决定作业能不能稳定跑。
4.1 RocksDB 还是 HashMap:看状态量和恢复时间
HashMapStateBackend 把状态放 JVM 堆内存,读写快,但状态量受堆大小限制,而且 Checkpoint 时要把整个状态序列化,状态大了之后 Checkpoint 时间线性增长。RocksDBStateBackend 把状态放本地磁盘,支持增量 Checkpoint,状态几百 GB 也能扛,代价是读写有序列化开销,延迟比 HashMap 高一些。
保险场景我一般选 RocksDB,原因是状态量会随业务增长,HashMap 迟早撑不住,中途换状态后端要从 savepoint 恢复,风险大。RocksDB 的增量 Checkpoint 在状态量大时优势明显,第一次全量之后,后续每次只传变更的 SST 文件。
# flink-conf.yaml 关键配置 state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: hdfs:///flink/checkpoints/insurance-risk state.savepoints.dir: hdfs:///flink/savepoints/insurance-risk state.backend.rocksdb.localdir: /data/flink/rocksdb state.backend.rocksdb.memory.managed: true state.backend.rocksdb.memory.write-buffer-ratio: 0.5 state.backend.rocksdb.memory.high-prio-pool-ratio: 0.1state.backend.rocksdb.memory.managed: true让 RocksDB 的内存由 Flink 的 MemoryManager 统一管理,避免和 JVM 堆争内存。write-buffer-ratio设 0.5 表示一半的 managed memory 给写缓冲,写多读少的场景可以调高。localdir要指向 SSD,机械盘上 RocksDB 的随机读写会成为瓶颈。
4.2 Checkpoint 间隔与超时:别让 Checkpoint 拖垮作业
Checkpoint 间隔太短,频繁触发影响吞吐;太长,故障恢复时回放的数据多。保险场景一般设 3 到 5 分钟。超时时间要大于单次 Checkpoint 的最长耗时,否则会频繁超时失败。
execution.checkpointing.interval: 3min execution.checkpointing.timeout: 10min execution.checkpointing.min-pause: 1min execution.checkpointing.max-concurrent-checkpoints: 1 execution.checkpointing.tolerable-failed-checkpoints: 3 execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATIONmax-concurrent-checkpoints设 1 是保险场景的保守选择,避免多个 Checkpoint 同时写 HDFS 抢带宽。tolerable-failed-checkpoints设 3 表示连续失败 3 次才让作业失败,给临时网络抖动留缓冲。externalized-checkpoint-retention设RETAIN_ON_CANCELLATION,手动取消作业时保留 Checkpoint,方便排查问题。
4.3 反压排查:火焰图怎么看
作业跑着跑着变慢,最常见的原因是反压。Flink Web UI 的 Backpressure 页面能看到每个算子的反压状态,但定位到具体代码要靠火焰图。开启火焰图需要在flink-conf.yaml里配rest.flamegraph.enabled: true,然后在 Web UI 的 Flame Graph 页面采样。
看火焰图时重点看哪个方法占的 CPU 时间最长。如果是RocksDB.get占比高,说明状态读多,考虑加缓存或优化 key 设计;如果是Serialization占比高,说明序列化开销大,考虑用更紧凑的数据类型;如果是某个业务方法占比高,那就是逻辑本身慢,要优化算法。我遇到过一次反压,火焰图显示BigDecimal.add占了 40% CPU,原因是理赔金额累加用了 BigDecimal,后来改成 long 存分,性能提升明显。
5. 避坑与排查:保险 Flink 作业最容易翻车的五个地方
5.1 保单状态乱序导致规则误判
现象:同一份保单的「退保」事件先到,「投保」事件后到,规则判断时发现「未投保先退保」,触发异常告警。
原因:CDC 抓 binlog 时,不同表的 binlog 到达 Kafka 的顺序不保证;或者 MQ 发送端用了异步发送,消息乱序。
解决:在 Flink 里用 EventTime 加 Watermark,设置合理的outOfOrderness,比如 10 秒。对于同一policy_id的事件,用keyBy(policy_id)保证同一保单的事件进同一个子任务,再在KeyedProcessFunction里按event_time排序缓存。如果乱序严重,考虑在 source 端做一次排序,或者用keyBy加ProcessFunction做本地排序。
5.2 Checkpoint 一直失败但作业不报错
现象:Web UI 上 Checkpoint 历史全是失败,但作业照常运行,直到某天重启时发现没有可用的 Checkpoint。
原因:tolerable-failed-checkpoints设得太大,或者 Checkpoint 失败只打 WARN 日志没引起注意。
解决:把tolerable-failed-checkpoints设小一点,比如 3;同时配监控,Checkpoint 连续失败 2 次就告警。检查 HDFS 权限和空间,Checkpoint 目录写不进去是最常见的原因。另外 RocksDB 的localdir磁盘满了也会导致 Checkpoint 失败,要监控磁盘使用率。
5.3 状态无限增长把 RocksDB 撑爆
现象:作业跑了一周之后,RocksDB 本地磁盘占用从 10GB 涨到 200GB,TaskManager 频繁 Full GC。
原因:ListState或MapState里的数据没有清理逻辑,或者定时器注册了但onTimer里没做清理。
解决:每个状态都要有对应的清理策略。ListState用定时器定期清理过期项;MapState设 TTL,用StateTtlConfig配置。TTL 的StateTtlConfig.newBuilder(Time.days(7)).setUpdateType(OnCreateAndWrite).setStateVisibility(NeverReturnExpired).build(),注意NeverReturnExpired表示过期数据不返回但也不立即删除,实际删除要等 RocksDB compaction。
5.4 广播规则更新后部分子任务没生效
现象:发了新规则到rule_topic,但部分被保人的告警还是按旧规则触发。
原因:BroadcastState 的更新是异步的,processBroadcastElement在每个子任务里独立执行,如果某个子任务的广播流消费滞后,规则就没更新。
解决:广播流的并行度要和主流程一致,且广播流的 source 要保证所有子任务都能消费到全量规则。常见做法是广播流用parallelism = 1的 source 读规则,然后broadcast()会自动把规则发给所有下游子任务。另外规则更新后可以发一条「版本号」消息,主流程收到事件时检查规则版本,版本不一致时等待或告警。
5.5 理赔高峰期 Kafka 消费滞后
现象:白天理赔高峰期,Kafka 消费 lag 从 0 涨到几十万,告警延迟从秒级变成分钟级。
原因:Flink 作业的并行度不够,或者单个子任务处理逻辑太重(比如同步调用外部风控接口)。
解决:先看 Kafka 分区数,Flink source 的并行度不能超过分区数,否则有子任务空闲。如果分区数够但并行度不够,加并行度。如果是逻辑重,把同步调用改成异步 IO,用AsyncDataStream.unorderedWait并发调用外部接口。保险场景里核保规则经常要查外部黑名单,同步查一次几百毫秒,异步化之后吞吐能翻几倍。
6. 用 Savepoint 做规则灰度与回滚:一个保险上线的具体技巧
保险作业上线最怕的是新规则误杀正常保单。我一般会用 Savepoint 加双跑做灰度:新规则先在一个低并行度的作业里跑,只输出告警不拦截,和旧作业的告警做对比,确认新规则没有大量误报之后再切流量。
具体操作分三步。第一步,从当前作业触发 Savepoint:
flink savepoint <job-id> hdfs:///flink/savepoints/insurance-risk-v1第二步,用 Savepoint 启动新作业,新作业的规则配置指向新规则集,但输出到一个单独的 Kafka topic 做观察:
flink run -s hdfs:///flink/savepoints/insurance-risk-v1 \ -c com.insurance.RiskJob \ /opt/flink/jobs/insurance-risk-v2.jar \ --rules.topic rule_topic_v2 \ --alert.topic alert_topic_shadow第三步,对比两个 topic 的告警。如果新规则的告警量和旧规则差异在可接受范围内(比如误报率不超过 1%),再把新作业的alert.topic切到正式 topic,停掉旧作业。如果发现新规则有问题,直接从 Savepoint 恢复旧作业,回滚时间取决于 Savepoint 大小,几百 GB 的状态大概几分钟。
这个流程的关键是 Savepoint 要包含所有状态,所以作业里不能有transient之外的非托管状态。另外 Savepoint 的路径要有版本管理,每次上线前打一个 tag,方便回溯。我习惯在 Savepoint 路径里带上日期和版本号,比如insurance-risk-20250115-v2,半年后回头看还能找到当时的现场。
还有一个细节:Savepoint 恢复时如果作业的算子 UID 变了,恢复会失败。所以每个算子都要显式设uid(),比如keyBy(...).process(...).uid("underwrite-rule")。这个 UID 一旦上线就不能改,改了就得从零开始跑,状态全丢。这是我在保险项目里踩过的最贵的坑,一次改 UID 导致状态全丢,重新积累状态花了两周。
希望帮到你。
本文还有配套的精品资源,点击获取