简介:面向保险行业大数据与实时数仓方向的 Spark 实战资源包,适合具备一定 Scala 与 Hadoop 基础、希望掌握 Kafka+Spark Streaming+HBase 整条实时链路的学习者。项目围绕业务系统数据库实时同步与统计报表展开,完整演示数据采集、流式计算、结果入库到报表展示的闭环,解决业务库变更事件难以实时转化为可读指标的痛点。压缩包共 400 个文件、大小仅 632KB,以 28 个 Scala 源文件为主体,配合 Sample 样例数据、Shell 启动脚本、Properties/XML 配置等,正好覆盖从环境配置到任务提交的各个调试环节。已吸引 979 人学习,说明该案例在保险实时数据分析场景中具有较高的参考价值。通过学习可拿到一套可运行的实时统计项目骨架,包括数据消费逻辑、窗口聚合算子、HBase 读写实现与常用部署脚本,适合作为真实业务项目的起步模板。
1. 保险行业跑 Spark 差在哪:一张保单背后的多源数据账
做保险数仓的人第一次把 Spark 跑在全量保单上,通常会得出一个反直觉的结论:瓶颈不在计算。几千万张保单、几亿条理赔流水,Spark 半小时能算完;真正吃掉一整天的是数据问题——同一张保单在不同系统里日期格式不一样,理赔金额在明细和汇总表里口径不一致,代理人维度的数据倾斜能把一个 join 从三分钟拖到三小时。这个「Spark 实战项目(保险行业真实项目)」要解决的就是这件事:把保单、理赔、续期这些多源异构数据搬进 Spark,按客户和产品维度算清楚赔付率、续期率和风险分层。适合正在做保险数据平台、数仓或报表开发的人,也适合想找一个结构化业务练手的数据工程师。
2. 把保险数据搬进 Spark:多源接入、schema 与会话参数
2.1 保险数据的五个来源,先建 ODS 层再谈口径
保险公司的数据源比互联网公司更杂。核心承保系统一般是 Oracle 或 SQL Server,里面是保单主表、批改记录;理赔系统独立一套,立案、医疗清单、赔付明细各自成表;收付费系统管实收保费和佣金;第三方渠道(经代机构、网销平台)按月导 CSV 或 JSON 订单日志;再加上精算和财务手工补录的 Excel 转件。这五类数据第一次拉通时,最典型的翻车是:拿一个 CSV 全量读入,发现时间字段有的被识别成字符串、有的识别成时间戳,地区码有的是两位、有的是四位,身份证号被读成 double 后精度直接丢。
做这个项目,我建议先建 ODS 层,只做格式标准化,不做业务口径合并。ODS 层每个表先定两件事:主键和粒度。保单表的粒度是一张保单一个 policy_no,理赔表的粒度是一次立案一个 claim_no,理赔明细表的粒度是一张医疗清单一个 id。粒度一旦定错,后面所有 join 和聚合都建立在流沙上——理赔金额重复统计(见 5.3)就是粒度混乱的典型后果。这一步不需要 Spark 跑得多快,但要保证每一张表下游能回答「一行代表什么」这个问题。
2.2 读保单 CSV:显式 schema 与 badRecordsPath 才是正解
保险行业导出的 CSV 很少是标准逗号分隔,字段里带逗号是常事,所以大家默认用竖线|或制表符。读取时最忌讳让 Spark 自己推断类型,原因很实际:18 位身份证会被推断成 long 然后荡失精度变成科学计数法;金额列如果混入空值和千分位逗号,整列被推断成 string;日期列在不同月份导出时格式可能漂移。下面这个写法是项目里的标准开头。
import org.apache.spark.sql.types._ // 保单主表:列很多,先只挑口径相关字段 val policySchema = StructType(Array( StructField("policy_no", StringType, true), // 保单号,字符串,不能当数值 StructField("holder_id_card", StringType, true), // 身份证号 18 位,绝不读成 Long StructField("holder_name", StringType, true), StructField("product_code", StringType, true), StructField("premium", DoubleType, true), // 承保保费,单位:元 StructField("policy_status", StringType, true), // 有效 / 退保 / 满期 StructField("start_date", DateType, true), StructField("end_date", DateType, true) )) val policyDF = spark.read .option("header", "true") .option("delimiter", "|") .option("dateFormat", "yyyy-MM-dd") .option("badRecordsPath", "/data/insur/bad/policy_202401/") .schema(policySchema) .csv("/data/insur/ods/policy/202401/")这个写法里有三个参数是保险场景的救命稻草。delimiter指定竖线分隔,避免地址、备注字段里的逗号把列数打乱;dateFormat在读取源头就统一日期,比后面到处to_date再处理 null 干净得多;badRecordsPath让 Spark 把解析失败的坏行单独落盘,而不是用DROPMALFORMED悄悄丢数据。保险数据是要过审计的,「丢了一行」在业务上说不清,把坏行留到一张异常表里,至少能回答「哪一行因为什么没进来」。
2.3 从核心系统 JDBC 拉数:分区列与 fetchsize 的取舍
保险公司的核心库通常只给只读账号,也不会允许你在生产库上跑复杂聚合。常见做法是用 Spark 的 JDBC 数据源并行拉取保单主表,落到 Hive 之后再算。这里有几个参数直接决定拉数速度,写错一个都可能把源库连接池打满或者慢到怀疑人生。
val policyFromJdbc = spark.read .format("jdbc") .option("url", "jdbc:oracle:thin:@//10.20.1.5:1521/ORCLPDB") .option("dbtable", "(select policy_no, holder_id, premium, start_date from t_policy where settle_date >= '2024-01-01') t") .option("user", "spark_etl") .option("password", "******") .option("fetchsize", "5000") .option("partitionColumn", "policy_id") .option("lowerBound", 1) .option("upperBound", 300000000) .option("numPartitions", 12) .load()partitionColumn必须选数值列(常见用自增主键),Spark 会把lowerBound到upperBound均匀切成numPartitions段分别去查,把一次大查询拆成 12 个并行查询。注意这个字段不能是日期,日期列做不了范围切分,硬写会扫全表。fetchsize一定要设,Oracle 驱动默认每次网络往返只拉 10 行,几百万行数据在这种配置下能跑到天荒地老,设到 3000~5000 是实践中的甜点区间。numPartitions不要贪大,它受源库允许的并发会话数和连接池大小限制,保险核心库通常很敏感,12 到 20 就够了,不是越大越快。
3. 用 DataFrame 算清客户风险账:赔付率、聚合与窗口函数
3.1 赔付率口径:已赚保费和已发生赔款为什么不能直接相除
如果只是跑技术 demo,赔付率就是一个除法;但在保险业务里,这个数直接进经营分析会和精算打架。第一年收的多年期保费不能全部算进当期收入,要按责任期分摊成已赚保费;赔款也分已决和未决,未决赔款对应的是精算计提的准备金,不能混进去。所以代码里的sum(premium)必须先明确口径:是承保保费还是已赚保费。下面这个写法是按评估期与保障期的重叠天数做分摊的简化版,做经营月报时最常用。
import org.apache.spark.sql.functions._ // 评估期:2024-01-01 到 2024-12-31(已赚保费口径) val evalStart = java.sql.Date.valueOf("2024-01-01") val evalEnd = java.sql.Date.valueOf("2024-12-31") val premiumDF = policyDF .withColumn("cov_start", greatest(col("start_date"), lit(evalStart))) // 保障起点和评估起点取晚者 .withColumn("cov_end", least(col("end_date"), lit(evalEnd))) // 保障终点和评估终点取早者 .withColumn("cov_days", datediff(col("cov_end"), col("cov_start"))) // 重叠天数 .withColumn("earned_premium", col("premium") * col("cov_days") / 365.0) // 一年期按天数分摊这个逻辑有三个边界坑。跨年保单的重叠天数容易差一天,金融口径会要求算头不算尾或反着来,跟财务确认一次再定死;退保保单的end_date被批改成了退保日,要比原终止日早,所以least里的选择很重要;短期险(航意险、旅游险这种保障几天的)完全不适合按 365 天分摊,要按保障期实际天数做分母,项目里一般单独建一张「短期险分摊例外表」处理。做这类字段,我一般不会直接写死在主流程里,而是先产出一张中间表,让业务核对过数值再进下游。
3.2 客户维度聚合:left join 保留无赔客户,groupBy 前想好粒度
客户维度是保险经营分析的地基。一个客户可能买了多张保单、出过多次险、退过保、又续过费,要把这些行为压缩成一行指标,标准做法是保单表左连接理赔表再按客户分组。下面这一段是整个项目的核心聚合之一。
// 客户维度风险画像:把保单和理赔关联后按客户分组 val customerRisk = policyDF .join(claimDF, Seq("policy_no"), "left") // 保留无赔客户,left 而不是 inner .groupBy("holder_id_card", "holder_name") .agg( sum("earned_premium").as("total_earned_premium"), count("claim_no").as("claim_count"), // 无赔客户 count 出来是 0,正好符合预期 sum("claim_amount").as("total_claim_amount") // 已决赔款口径,未决另算 )这里选left join不是拍脑袋:续期率和客户分层要用到全部客户,只 join 出过险的人,指标就废了。count("claim_no")在左连接后对无赔客户返回 0,这是想要的语义;如果误用count(*)会把无赔客户算成 1 次出险,这种错很隐蔽。sum对 null 自然返回 null,但有些报表工具不认 null,建议聚合后统一coalesce(..., 0)。另外要清楚 shuffle 发生在groupBy这一行:相同客户的所有数据会被拉到同一个 executor 上做最终聚合,数据倾斜最容易发生在这里,具体解法放到第 5 章。
3.3 窗口函数算续期:用 lag 看客户断档天数
客户维度的第二个常用指标是续期行为:这个客户上一张保单到期之后多少天买了下一张。这个用窗口函数写最直接,比自连接干净很多。
import org.apache.spark.sql.expressions.Window val w = Window.partitionBy("holder_id_card").orderBy("start_date") val renewalDF = policyDF .withColumn("prev_end", lag("end_date", 1).over(w)) // 上一张保单的到期日 .withColumn("gap_days", datediff(col("start_date"), col("prev_end"))) // 断档天数lag取同一个客户按保单生效日排序的前一行end_date,算出来的gap_days可以打标签:小于 0 说明新旧保单责任期重叠,0 到 30 天算常规续保,大于 30 天算犹豫后回归,null 则是新客首单。窗口函数和groupBy的区别在于它不折叠行数,每一张保单保留自己的位置,这在做「第几单续保」这类序列分析时是必要的。代价是窗口函数要在分区内保存排序数据,客户量大时内存会被吃掉不少,spill 到磁盘会拖慢速度;如果只是算客户级汇总,还是优先用 3.2 的groupBy,窗口函数留给真正需要看序列的场景。
4. 让 Spark 作业稳定跑在保险集群上:参数、UI 与 AQE
4.1 三个必调参数:executor 内存、分区数与并行度
保险行业的 Spark 作业大多是定时跑批,不像互联网场景有持续的实时流量,所以参数配置的核心目标是「稳定跑完、不抢资源、可重跑」。下面这套 spark-submit 参数是项目里调过很多轮之后沉淀下来的基准配置。
spark-submit \ --class com.insur.risk.CustomerRiskJob \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.sql.adaptive.enabled=true \ --conf spark.sql.adaptive.skewJoin.enabled=true \ --conf spark.driver.maxResultSize=2g \ /data/insur/jars/customer-risk.jar先说executor-memory。8g 指的是堆内内存,而 YARN 在申请容器时还会额外算一块 overhead(默认是 max(384MB, 0.1 × executor 内存)),所以实际容器内存大约是 8g + 819MB。如果物理机只有 16g,20 个 executor 加 4g driver 会直接把队列资源打满,其他团队的任务全被挤到等待。保险团队一般不会只有一个跑批任务,资源要留给别人,num-executors建议按「总数据量 ÷ 单 executor 处理能力」来算,而不是盲目乘大。再就是spark.sql.shuffle.partitions,默认 200 对几亿行理赔明细完全不够,每个分区几百万行会让单个 task 处理时间过长;我一般按「预期输出大小 ÷ 128MB」估算分区数,再留 20% 余量。
spark.driver.maxResultSize是另一个容易忽略的点:如果代码里有collect()把聚合结果拉回 driver,不设上限会直接把 driver 内存撑爆。跑批项目里我基本禁止collect大结果,只允许拉取几十行的汇总到 driver 打印日志。
4.2 用 Spark UI 定位内存瓶颈:GC 与 shuffle spill 怎么看
任务卡死了,第一反应不是加内存,而是打开 Spark UI。做 spark 内存线程监测工具选型时,大家最后都会回到 UI 自带的那几张页面——够用且不用额外部署。重点看两处:Stage 页的 Summary Metrics 里有没有Shuffle Spill (Disk),以及 Executors 页每个 executor 的Shuffle Read/Write是否均匀。
我遇到过一次典型的客户聚合卡死:1000 个 task 里 999 个 2 秒跑完,剩一个跑了 2000 秒还不出结果。打开 UI 看,那个 task 的 input 数据量是中位数的 15 倍,GC 时间占了 task 时间的 40%——这不是内存不够,是数据倾斜。加内存只会让那个倒霉的 executor 多撑一会儿,治标不治本。另一次是某个作业频繁报Container killed by YARN,撑开 executor 日志看是堆外内存超限,最后把spark.executor.memoryOverhead从默认值调大才解决。所以看 UI 的顺序我固定成这样:先看有没有 spill,再看 GC 占比,再看各 task 的 input 分布,三样都正常才谈加资源。
4.3 AQE 该开哪几个:skewJoin 与 coalesce 的取舍
Spark 3.x 的 AQE(自适应查询执行)对保险这类跑批作业是实打实的收益,因为跑批数据量大、倾斜又常见。spark.sql.adaptive.enabled是总开关,spark.sql.adaptive.skewJoin.enabled会在 join 阶段自动把大 key 拆成多份并行处理,spark.sql.adaptive.coalescePartitions.enabled会在 shuffle 后把过小的分区合并掉,减少空转的 task。
但要清楚 AQE 的边界。skewJoin 只对 join 阶段的倾斜有效,groupBy之后的倾斜它管不了,那还是要靠两阶段聚合或者加盐(见 5.1)。coalesce 分区合并对小输出很友好,但如果目标表本身需要按日期分区写 200 个分区,合并反而会让写入并行度下降。我的习惯是:总开关打开,skewJoin 打开,coalesce 在跑批环节关掉,等最后写表的算子前再repartition成目标分区数。AQE 不是配了就不看 UI,它只是把一部分调参交给了运行时的统计信息,排错思路一点没变。
5. 避坑指南:保险 Spark 项目的 5 个翻车现场
5.1 数据倾斜:代理人保单量悬殊,一个 task 拖垮整个 stage
现象:客户维度 join 跑了几分钟,其他 task 都结束了,只有一个或两个 task 还在跑,集群资源看着有一半是空的。
原因:保险销售的二八效应非常严重。头部代理人的保单量可能是普通代理人的几百倍,按agent_code或holder_id做 shuffle 时,某几个 key 对应的数据被哈希到同一个分区,那个 task 要处理的数据量级和其他 task 差两个数量级。
解决:老办法是加盐。把大表的小 key 加一个随机前缀,小表把同一个 key 复制 10 份再交配,把大 key 的负载匀到 10 个分区里。Spark 3 的 AQE 开了skewJoin之后能自动处理部分场景,但 groupBy 后的倾斜还是要手动打两阶段聚合:先加盐聚合一次,再去掉盐做最终聚合。注意加盐不能用在需要精确唯一性校验的场景,否则去重逻辑会被破坏。
import org.apache.spark.sql.functions._ // 加盐:把小表按盐值复制 10 份,大表随机打散,join 时均摊大 key val saltRange = (0 until 10).toArray val explodedClaim = claimDF .withColumn("salt", explode(lit(saltRange))) // 复制 10 份 val saltedPolicy = policyDF .withColumn("salt", when(col("agent_code").isNotNull, floor(rand() * 10).cast("int")) .otherwise(0)) val joined = saltedPolicy.join(explodedClaim, Seq("policy_no", "salt"), "left")5.2 日期字段混存两种格式,过滤结果悄悄变少
现象:某个月赔付率报表数字明显偏低,查了业务系统觉得没问题,但 Spark 算出来的理赔数据比系统少了几千条。
原因:核心系统历年导出日期格式不一致,老数据是yyyyMMdd,新数据是yyyy-MM-dd,甚至同一张表里两个格式都有。to_date解析失败时返回 null,后续where过滤把 null 全部滤掉了,数据悄悄变少而且没有任何报错。
解决:读取时用显式 dateFormat 只能管住一种格式,混存时最可靠的是先统一清洗:用一个 UDF 同时尝试几种格式,解析失败的整行进异常表,不让它静默消失。这个异常表要定期让数仓管理员看,因为它的存在本身就说明上游数据质量有问题。给下游的作业最好再加一条校验:聚合结果的行数若比前一天少超过阈值,直接终止任务并告警。
5.3 理赔金额重复统计:1 对多 join 后 sum 翻倍
现象:按产品汇总的赔款金额比财务对账结果多了 30%,业务追问下来发现是某类保单被重复算了好几次。
原因:保单表和理赔表是 1 对多,理赔表和理赔明细表又是 1 对多,连续两次 join 之后行数膨胀,sum(claim_amount)把同一笔赔款在明细行上累加了一遍。最典型的是「一张保单对应一次立案,一次立案对应五张医疗清单」,join 完五张清单后赔款金额被算了五遍。
解决:先按最细粒度聚合再去 join,而不是先 join 再聚合。先把理赔明细按 policy_no 把金额 sum 成一行,再和保单关联,从根上避免膨胀。
// 先按 policy_no 聚合理赔金额,再与保单表关联 val claimAgg = claimDetailDF .groupBy("policy_no") .agg(sum("paid_amount").as("paid_amount_total")) val resultDF = policyDF .join(claimAgg, Seq("policy_no"), "left")5.4 按天分区写入,小文件把 HDFS 元数据打爆
现象:任务跑完了,却发现 HDFS 上某个分区目录下有几千上万个小文件,每个只有几百 KB。下游再读这张表时扫描耗时成倍增长,NameNode 的元数据压力也跟着上来。
原因:spark.sql.shuffle.partitions设大了之后,每个 task 都会往目标分区写一个文件,200 个分区任务对 30 个日期分区,每个分区下平均产生 6~7 个文件,整个表就是几百个文件。日积月累,文件数爆炸。
解决:写表前先repartition到目标分区数的整数倍,而不是依赖 shuffle 后的分区数。更稳妥的做法是每分区只写一个文件:结果数据量不大时,repartition(分区数)最直接;数据量大时按分区键repartitionByRange或写完后跑一次合并小文件的作业。注意coalesce是窄依赖可能不会真正合并,保险起见用repartition。日常跑批调小spark.sql.shuffle.partitions或对写表单独开一个spark-submit任务,是两种常见方案,我偏向后者——避免拖慢前面的聚合步骤。
5.5 动态分区 OOM:不是内存不够,是分区数超出预估
现象:开启动态分区写 Hive 表,作业跑了一半直接 OOM,看日志是 driver 或 executor 内存溢出,但数据量并不大。
原因:动态分区写入时,每个 executor 要为它写入的每个分区各维护一个文件写入句柄。如果某个 task 接收的数据覆盖了上百个分区,句柄数乘上缓冲内存,很快就把 executor 堆内存打爆。保险数据按客户维度聚合后,客户可能分布在所有地区码和产品码上,分区组合一多就触发。
解决:预估分区数量,控制单批次写入的分区规模,或者干脆降低写入并行度,让每个 task 接手更少的分区。对日增量跑批,可以改成先写临时表、再按天INSERT OVERWRITE的方式,把动态分区拆成小批次。纯粹的调大spark.executor.memory只能推迟爆掉的时间,解决不了句柄数线性增长的问题。
6. 让结果被业务用起来:对账宽表与落表习惯
6.1 对账宽表:跟核心系统总账核对,别信自己的 sum
聚合结果跑完之后,第一件事不是画报表,而是对账。把 Spark 算出来的总保费、总赔款、保单数拿去和核心系统导出的总账比对,差异在万分之一以内才算通过。对不上就回到 5.3 检查 join 链路,或者回到 5.2 看是不是有日期解析失败被过滤。下面这段 SQL 是我每次跑完必须执行的。
-- 对账:Spark 结果对比核心系统总账,差异超过万分之一需回溯 select sum(premium) as spark_premium, sum(claim_amount) as spark_claim, count(distinct policy_no) as spark_policy_cnt from dws_risk_customer where stat_date = '2024-12-31';这个习惯救过我很多次。有一回差异出在千分之五,追了两天发现是某个产品线的退保保单policy_status更新不及时,批改系统还没把状态推给数仓。这种问题如果不做总量勾稽,散落到客户维度报表里根本看不出来。
6.2 明细表、宽表、异常表分开落,重跑只 drop 分区
最终落地我一般分成三类表。明细粒度结果表给风控模型做抽样校验用,粒度最细、字段最少;客户宽表给报表工具做经营分析,字段也不宜太宽,超过 50 个列的宽表在保险报表场景里查询性能很差;异常数据表给数仓管理员,记录所有被清洗规则拦下来或解析失败的数据,方便他们回头找上游系统要数。这三类表用同一个批处理脚本串起来:先跑清洗,再跑聚合,最后对账,任一步失败就终止,不往 Hive 里写半成品。
重跑逻辑上只 drop 当天分区再写入,不 drop 整张表,保留历史分区以便回滚查数。这套流程补齐了 Spark 作业之外的工程习惯,也才让「Spark 实战项目」真正从跑通变成了能交付给业务用的数仓作业链。现在我已经把「先对账、再看异常表、最后看指标表」固化成肌肉记忆了——数据落表之前谁都不敢说结果是对的。希望帮到你。
本文还有配套的精品资源,点击获取