保险行业Spark数仓实战:多源接入、数据倾斜与赔付率计算

发布时间:2026/10/3 10:33:36
保险行业Spark数仓实战:多源接入、数据倾斜与赔付率计算
简介面向保险行业大数据与实时数仓方向的 Spark 实战资源包适合具备一定 Scala 与 Hadoop 基础、希望掌握 KafkaSpark StreamingHBase 整条实时链路的学习者。项目围绕业务系统数据库实时同步与统计报表展开完整演示数据采集、流式计算、结果入库到报表展示的闭环解决业务库变更事件难以实时转化为可读指标的痛点。压缩包共 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 行几百万行数据在这种配置下能跑到天荒地老设到 30005000 是实践中的甜点区间。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 与 AQE4.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.partitions200 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.skewJoin.enabledtrue \ --conf spark.driver.maxResultSize2g \ /data/insur/jars/customer-risk.jar先说executor-memory。8g 指的是堆内内存而 YARN 在申请容器时还会额外算一块 overhead默认是 max(384MB, 0.1 × executor 内存)所以实际容器内存大约是 8g 819MB。如果物理机只有 16g20 个 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 个日期分区每个分区下平均产生 67 个文件整个表就是几百个文件。日积月累文件数爆炸。解决写表前先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 实战项目」真正从跑通变成了能交付给业务用的数仓作业链。现在我已经把「先对账、再看异常表、最后看指标表」固化成肌肉记忆了——数据落表之前谁都不敢说结果是对的。希望帮到你。本文还有配套的精品资源点击获取