首页 / 资讯中心 / 文章详情

Spark保险理赔风险预测项目实战:从数据清洗到模型上线

Spark保险理赔风险预测项目实战:从数据清洗到模型上线 ★ FEATURED ARTICLE
简介面向大数据开发者与 Spark 学习者的保险行业实时数据分析实战项目资料包覆盖 Kafka 消息同步、Spark Streaming 实时计算、HBase 结果存储的完整链路解决业务数据库实时同步与统计报表分析场景适合具备一定 Scala 与大数据基础、希望动手实践流处理的读者。压缩包共400个文件含 Scala 源码、Shell 部署脚本、properties/XML 配置及 sample 样例数据各类型分别对应计算逻辑、启动部署、参数配置和输入样例整体仅632KB结构紧凑适合快速搭建实验环境。已有979人学习下载。其中的核心计算逻辑、组件对接方式、启动与配置脚本及样例输入配合实现流程中的窗口聚合、过滤等实时操作有助于还原 Kafka 数据接入、Spark 窗口计算与 HBase 写入的每个环节快速理解保险保费收入、理赔监控等实时报表的落地思路为进一步改造或迁移到生产环境提供参考。1. Spark实战项目保险行业真实项目到底在解决什么Spark实战项目保险行业真实项目听起来像一份课程作业但真正处理过保险数据的人知道它背后是一整套能落地的工程链路。车险公司每天进来上百万条报案、查勘、定损和结案记录单是清洗脏数据这一个环节用 Pandas 就能把一台 16G 内存的机器跑死。Spark 把这些数据分散到多台机器上并行处理再在这个基础上建特征、训练理赔风险评分模型。这个项目要解决的核心问题是帮保险公司从海量理赔数据中找出高风险案件并支撑精算、核保和反欺诈部门做决策。适合正在做保险数据平台、反欺诈系统或者想系统掌握 Spark 全流程的工程师。你不需要先搭一个数据中台只要有一台能提交 Spark 任务的环境就能把这条路走通。2. 从业务到表理赔风险预测项目的计算架构和数据模型2.1 为什么这个项目用 Spark 而不用 Pandas数据规模和计算模型决定选型保险公司的核心理赔明细表常见量级是几亿行字段数在 50 到 150 之间。Pandas 更适合单机交互式分析数据一旦超过内存就会触发 swap一个 groupby 能等半小时。Spark 的赢点在于把数据切成分区partition由多个 executor 并行跑DAG 调度器自动优化执行计划。从开发效率看Spark SQL 把 MapReduce 的门槛降得很低一个 JOIN 就是一行 SQL。更重要的一点是后续模型训练可以直接用 Spark MLlib不用把数据切块导到 Python 里再做一遍特征工程。整个流程从原始文件到模型输出全部留在 Spark 上省掉了大量数据搬运的 IO 开销。那是不是所有保险项目都该上 Spark不一定。如果你的业务就是几百万行维表每天做一次聚合分析单机 Pandas 更快没必要引入集群运维成本。只有当数据量超过单机峰值或者需要每天定时跑批处理、任务能容忍分钟级延迟时Spark 才是划算的选型。真实项目里我一般先看数据规模单表低于 5000 万行且列少直接 Pandas 做基线超过这个量级直接上 Spark省得后期换框架重来。2.2 保险业务数据怎么落到表理赔明细、保单、客户三张核心表一个典型的车险理赔风险项目至少要准备三类数据。理赔明细表claim是主表包含报案号、保单号、出险时间、报案时间、事故类型、理赔金额、责任比例、结案状态。保单表policy提供被保人年龄、车辆价值、保额、投保险种组合。客户表customer记录历史出险次数、历史赔款金额、续保年限这些是构建风险画像的核心特征。三张表的关联关系很直接理赔表和保单表通过保单号关联保单表和客户表通过客户 ID 关联。落到数仓里我一般把它们按 parquet 格式存储Spark 读 parquet 比读 CSV 快一倍左右还自带压缩。表字段设计有个容易忽略的点日期字段既可能存成毫秒时间戳也可能存成“2024-03-15 10:30:00”或“2024/03/15 10:30:00”。这种脏格式会在后面清洗阶段统一处理。另外业务系统中的“责任比例”字段可能存成“全责”“50%”“0.5”三种形态模型不能直接用必须提前标准化。表名关键字段用途claimclaim_id, policy_no, accident_time, report_time, claim_amount, accident_responsibility理赔明细左表作为建模主表policypolicy_no, age, vehicle_value, coverage_code保单维表提供人和车的信息customercustomer_id, historical_claims, historical_loss_amount客户维表提供历史风险特征有些项目还会引入修车厂维表、地域维表这要看业务侧能不能拿到干净的数据源。我建议第一版先做三张表把主链路跑通再往外扩维一上来就接十几个表和几十个字段清洗负担会让你连基线模型都跑不出来。2.3 用 spark-submit 提交任务最小可运行的环境配置和初始化代码我习惯把整个项目写成一个 Python 文件然后用 spark-submit 提交到 YARN 集群。下面是一份最基础的提交命令spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ insurance_risk.pymaster yarn表示让 YARN 做资源调度deploy-mode cluster表示 Driver 也跑在集群里提交命令的机器不会成为单点。driver-memory 4g是 Driver 的堆内存executor-memory 8g是每个 Executor 的堆内存executor-cores 4是每个 Executor 占用的 CPU 核数num-executors 20是启用的 Executor 数量。这些数字必须结合集群剩余资源调整不能照抄如果是在单机伪分布式环境调试把--master改成local[*]去掉 YARN 相关配置即可。项目代码里的 SparkSession 初始化也有一点讲究from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(insurance_claim_risk) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate()appName用于在 YARN 和 Spark UI 里识别任务。spark.sql.shuffle.partitions控制所有 shuffle 操作JOIN、GROUP BY的默认分区数默认值是 200但它不是越大越好——分区数超过实际数据量会产生大量空任务反而拉长调度时间。enableHiveSupport()是当你要读写 Hive 表时打开否则可以去掉。这个初始化代码建议放在文件最开头后面所有的 DataFrame 操作都复用同一个 SparkSession。环境起来之后先用最基础的手段探查数据。读取理赔表的 parquet 文件看字段类型和行数df spark.read.parquet(/data/insurance/claim) df.printSchema() total df.count() df.groupBy(accident_type).count().show(10)printSchema()打印每个字段的类型count()是一个 action会触发全表扫描几亿行的表可能要跑几分钟。如果只想看事故类型分布groupBy加上show(10)就够。这里有个经验对于超大表我很少在开发阶段直接count()而是用approx_count_distinct或者先采样一个 1% 的子集快速验证逻辑再跑全量。不然每调一次代码就全表扫一次集群资源很快就抗不住。3. 保险理赔数据的清洗与标准化用 Spark SQL 把脏数据变干净3.1 缺失率高的列先统计再决定删还是填保险数据里“第三者责任险保额”这种字段经常有近一半为空原因是很多保单根本没买这个险种。对于这种高缺失列直接删掉比硬填充更安全因为它实际上代表一种业务状态填充反而引入噪声。另一类是理赔金额为空这通常意味着案件还没结案缺失比例不高可以用中位数填充。第一步永远是统计每一列的缺失率而不是凭感觉决定怎么处理。下面这段代码会输出每个字段的缺失率from pyspark.sql.functions import col, count total df.count() missing_stats df.select([ (1 - count(col(c)) / total).alias(c) for c in df.columns ]) missing_stats.show()count(col(c))计算的是非空行数除以总行数再取反就得到缺失率。注意df.count()在这里被调用了两次一次赋值一次在列表推导式里如果数据量大最好先把total缓存住。实际项目里我会先跑一个采样版本的缺失率比如df.sample(0.01)算一遍等确定列处理逻辑后再跑全量。根据统计结果决定删除和填充策略df_clean df.drop(third_party_liability_limit) median_amount df.approxQuantile(claim_amount, [0.5], 0.01)[0] df_clean df_clean.na.fill({claim_amount: median_amount})drop直接删除整列。na.fill接受一个字典指定每列需要的填充值。approxQuantile返回的是近似中位数因为精确中位数需要全局排序开销很大第三个参数0.01表示相对误差数字越小越精确但越慢。对于理赔金额这种长尾分布中位数比均值更抗异常值。如果你的表特别大连近似分位数都觉得慢那就直接用 0 填充对于风控模型来说理赔金额为 0 也会被当作一个特征模式来学习。3.2 重复报案和日期乱码三个容易踩的数据坑重复报案是保险数据里真实存在的脏数据问题。同一个报案号可能出现多次原因是线下录入重复、接口重试或者多系统同步冲突。去重时不能只看报案号因为不同险种的报案号可能撞号。我一般用“报案号 出险时间”两个字段联合去重。日期格式混乱是第二个常见问题。同一张表里既有毫秒时间戳也有“2024-03-15 10:30:00”和“2024/03/15 10:30:00”两种字符串。可以分别用两个格式解析再合并from pyspark.sql.functions import to_timestamp, lit, when df_dedup df_clean.dropDuplicates([claim_id, accident_time]) df_std_time df_dedup \ .withColumn(ts_a, to_timestamp(col(accident_time), yyyy-MM-dd HH:mm:ss)) \ .withColumn(ts_b, to_timestamp(col(accident_time), yyyy/MM/dd HH:mm:ss)) \ .withColumn(accident_ts, when(col(ts_a).isNotNull(), col(ts_a)).otherwise(col(ts_b))) \ .drop(accident_time, ts_a, ts_b)to_timestamp支持传入格式字符串生成一个timestamp类型的新列。我们用两个临时列分别解析两种格式再用when优先取转换成功的那一个。解析失败的字符串会变成NULL所以最后留下来的NULL就是真正的格式错误。第三个坑是责任比例字段。业务上存“全责”“50%”“0.5”三种形态模型只能吃数值。需要统一转成 0 到 1 的浮点数from pyspark.sql.functions import col, when, lit df_std df_std_time.withColumn( resp_ratio, when(col(responsibility).contains(%), col(responsibility).cast(double) / 100) .when(col(responsibility) 全责, lit(1.0)) .when(col(responsibility) 无责, lit(0.0)) .otherwise(col(responsibility).cast(double)) )contains(%)先把字符串形式的“50%”识别出来cast(double)会把“50%”转成 50再除以 100 变成 0.5。when支持多分支按顺序评估。最后otherwise兜底直接把已经是数字的字符串转成 double。这里要小心如果字段里还有“主责”“次责”之类的文本otherwise返回的可能是NULL建议先做一层where过滤把异常值单独拉出来看。3.3 从明细表到宽表JOIN 时怎么避免数据膨胀和错误关联模型训练需要把保单和客户维表的字段合并到理赔明细表上。常见做法是left join以理赔明细为主表。但维表如果很大JOIN 会产生大量 shuffle拖慢整个任务。一个实用技巧是先对维表做投影和去重只保留要用的列再把小维表广播出去。policy_dim spark.read.parquet(/data/insurance/policy) \ .select(policy_no, age, vehicle_value, coverage_code) \ .dropDuplicates([policy_no]) from pyspark.sql.functions import broadcast feature_df df_std.join(broadcast(policy_dim), policy_no, left)broadcast是一个提示函数告诉 Spark 把右表分发到每个 executor 内存中避免 shuffle join。适用条件一般是右表大小小于spark.sql.autoBroadcastJoinThreshold默认是 10MB。如果你的维表有几百 MB强行广播会让 driver 直接 OOM。这时候要么调大阈值要么改用分桶 join。真实项目中保单维表控制在 50MB 以内是比较健康的我一般会在导入阶段就压缩字段去掉大字段比如保单条款全文只留编码。还有个容易忽略的点如果一张保单在维表里有多条记录比如中途变更过保险内容join后理赔明细会成倍膨胀。所以在 join 前先做dropDuplicates是必需的。你可能觉得这多此一举但保险数据里保单变更是常态不加这一步等到训练时发现样本数量翻了三倍再回头查就费劲了。4. 特征工程与模型训练用 Spark MLlib 构建理赔风险评分模型4.1 特征处理把字符串和数值统一成向量输入MLlib 里的分类器吃的是稠密或稀疏特征向量。保险项目里常见特征类型有三类连续数值年龄、车辆价值、理赔金额、类别编码险种组合、车型代码、地域代码、计数特征历史出险次数。连续数值往往分布不均匀直接用原始值容易让树模型分裂点偏向高密度区间我习惯先用QuantileDiscretizer做分箱。类别特征需要先转成索引再做 one-hot。from pyspark.ml.feature import QuantileDiscretizer, OneHotEncoder, StringIndexer, VectorAssembler discretizer QuantileDiscretizer( numBuckets10, inputColage, outputColage_bucket ) indexer StringIndexer( inputColcoverage_code, outputColcoverage_code_idx, handleInvalidkeep ) encoder OneHotEncoder( inputColcoverage_code_idx, outputColcoverage_code_vec ) assembler VectorAssembler( inputCols[age_bucket, vehicle_value, historical_claims, coverage_code_vec], outputColfeatures )QuantileDiscretizer按数据分位数切割比等宽分箱更稳健numBuckets10意思是把年龄切到 10 个桶里。StringIndexer把类别字符串映射成数字索引handleInvalidkeep特别重要它会在训练时为新出现的类别预留一个特殊桶避免线上预测遇到未见过的险种编码直接报错。OneHotEncoder把索引展开成稀疏向量。VectorAssembler把分箱、one-hot 和数值拼成一个总特征向量这就是最终模型的输入。有一点需要注意StringIndexer在训练和预测时如果分开调用可能会因类别顺序不一致导致错位。所以后面一定要用Pipeline把整个处理链绑定起来而不是单独保存每个 transformer。4.2 训练模型用随机森林还是 GBDT理赔风险预测是个典型的二分类问题标签是“是否高风险案件”。正样本通常只占 5% 以内很不平衡。随机森林对不平衡数据相对稳能输出特征重要性也容易调参所以我一般先拿它做 baseline。Spark 里的GradientBoostedTrees拟合能力更强但调参时间长对偏斜分布更敏感。等随机森林基线跑通再决定要不要换。下面是训练和评估的核心代码from pyspark.ml.classification import RandomForestClassifier from pyspark.ml.evaluation import BinaryClassificationEvaluator train, test feature_df.randomSplit([0.8, 0.2], seed42) rf RandomForestClassifier( featuresColfeatures, labelColis_high_risk, numTrees100, maxDepth10, maxBins32, impuritygini, seed42 ) model rf.fit(train) evaluator BinaryClassificationEvaluator( labelColis_high_risk, metricNameareaUnderROC ) print(train AUC:, evaluator.evaluate(model.transform(train))) print(test AUC:, evaluator.evaluate(model.transform(test)))randomSplit([0.8, 0.2], seed42)的第二个参数是随机种子不要忽略没有 seed 的话每次跑出来的数据集划分都不同模型就会像开盲盒一样时好时坏。numTrees一般从 50 开始试100 是性能和稳定性的平衡点maxDepth默认是 5但理赔数据特征交互复杂10 层更实用再深容易过拟合。maxBins越大连续特征分裂点越精细代价是计算变慢32 是个常用起点。BinaryClassificationEvaluator的metricName选areaUnderROC打印出来的 AUC 在 0.7 到 0.8 之间算可用0.8 以上在保险反欺诈里已经不错。注意这里评估用的是全量测试集但没有处理样本不平衡所以 AUC 可能偏乐观后面第 5 章会专门聊这个问题。4.3 把模型持久化对全量明细打分训练完成后要做两件事保存模型然后用它跑全量数据打分。千万不能每次都重新训练那属于自找麻烦。最稳妥的做法是把所有预处理步骤和分类器装进一个Pipeline一次性保存。from pyspark.ml import Pipeline, PipelineModel pipeline Pipeline(stages[discretizer, indexer, encoder, assembler, rf]) pipeline_model pipeline.fit(train) pipeline_model.write().overwrite().save(/models/insurance_risk_pipeline) loaded_model PipelineModel.load(/models/insurance_risk_pipeline) full_df loaded_model.transform(feature_df) full_df.select(claim_id, probability, prediction).show(5)Pipeline里的每个 stage 会依次执行训练时只调用一次fit。保存到磁盘后线上任务直接PipelineModel.load加载。transform(feature_df)会对传入的 DataFrame 做完全一样的特征处理和预测。probability是一个向量第二个元素是正类概率业务系统通常把它乘以 1000 作为风险分来展示。保存模型到 HDFS 时有个细节overwrite()会先删除旧目录再写新的。如果你保存到同一个路径且没有加overwrite第二次跑会报文件已存在的异常。每次迭代模型时我建议在路径里带上日期比如/models/insurance_risk_pipeline_20250412既避免覆盖问题也方便回滚到上一版模型。5. 避坑Spark 跑保险大数据时最容易翻车的 5 个地方5.1 join 卡死在 reduce 阶段数据倾斜导致 shuffle OOM现象任务跑到某个 join 阶段进度停在 99% 超过十分钟然后容器报 OOM 被 kill。打开 Spark UI看到某个 partition 的 shuffle read 数据量是其他 partition 的几百倍。原因理赔数据里“经济型轿车”是最热门车型join 时同一个 key 的几百 GB 数据全部涌向同一个 executor。默认 hash 分区无法打散这种热点 key。解决给倾斜的 join key 加盐。思路是对热点 key 附加一个随机前缀把一份数据拆成多份并行处理join 完成后再去掉前缀还原。伪代码如下from pyspark.sql.functions import concat, lit, floor, rand # 对左表倾斜 key 加盐 df_left_salted df_left.withColumn( salted_key, concat(col(policy_no), lit(_), floor(rand() * 100)) ) # 右表先复制成 100 份保证能和每个加盐 key 对上 df_right_salted df_right.withColumn(range_id, explode(array(*(range(100))))) \ .withColumn(salted_key, concat(col(policy_no), lit(_), col(range_id))) df_joined df_left_salted.join(df_right_salted, salted_key)注意rand()不加种子即可加盐的目的只是打散热点不需要可复现。这个方案会增加 join 的中间数据量但对于几百 GB 的倾斜场景换来的是任务能成功跑完值得做。如果你用的是 Spark 3.4 以上版本可以先开启spark.sql.adaptive.skewJoin.enabledtrue让 Adaptive Query Execution 自动处理轻度倾斜只在阈值内。5.2 小文件多到下一阶段直接超长等待清洗任务写出的文件太碎现象每天清洗完理赔数据HDFS 上出现几千个只有几百 KB 的 parquet 小文件。下一次读这张表时任务启动和调度时间比实际计算时间还长。原因清洗流程里设置了过大的spark.sql.shuffle.partitions加上多次 action 后每个 task 都会写一个碎片文件最终生成的文件数量等于最后 stage 的 task 数。解决在写出之前用coalesce合并分区。注意coalesce(n)只减少分区不会触发布新 shuffle比repartition便宜很多。写出的 parquet 文件目标是单个文件在 128MB 左右这样读取效率最高。df_std_final df_std.coalesce(50) df_std_final.write.mode(overwrite).parquet(/data/insurance/claim_clean)coalesce(50)前面的数字要根据总数据量估算假设数据量约 6GB单文件 128MB那么 50 个分区是合理的。如果你的是几十亿行的表可以先跑一次approxCountDistinct估算行数再根据每行宽度估算总大小最后反推分区数。5.3 内存参数调了还是 OOMexecutor 内存不是越大越好现象Executor 内存从 8G 调到 16G 后任务仍然 OOM甚至更难跑完。原因Spark 内存分为执行内存execution和存储内存storage两者由spark.memory.fraction控制。内存调大后如果每个 partition 的 task 数不变单个 task 可能占更大内存反而加剧了 GC。另外如果数据倾斜没有被解决内存再大也有被单点数据打满的时候。解决遇到 OOM 不要只加内存先看 Spark UI 的 Executors 页面找到峰值内存最高的 executor确认是否数据倾斜。然后按顺序做三件事一是调大分区数spark.sql.shuffle.partitions或者spark.default.parallelism二是设置spark.memory.fraction0.75让执行内存占比更高三是开启spark.sql.adaptive.enabledtrue。这些参数用spark-submit --conf传递即可不用改代码。内存调优是有点玄学的部分我一般遵守“先分区后内存”的规律乱加内存多半是白花钱。5.4 日期字段解析后出现 8 小时偏差时区配置没有统一现象用to_timestamp解析“2024-03-15 10:30:00”后打印出来的时间变成了“02:30:00”。原因Spark 集群默认时区是 UTC而你的保险业务数据库用的是北京时间GMT8。to_timestamp解析字符串时不会自动带时区最终显示值取决于 SparkSession 的spark.session.timeZone配置。解决在初始化 SparkSession 时统一设置时区spark SparkSession.builder \ .appName(insurance_claim_risk) \ .config(spark.session.timeZone, Asia/Shanghai) \ .getOrCreate()设置后所有timestamp类型列的显示和计算都会基于北京时间。如果你不希望全局改时区也可以在查询时用from_utc_timestamp(col(accident_time), Asia/Shanghai)做转换。这个坑最容易出现在跨团队协作场景你本地用的默认时区是系统时区测试环境是 UTC线上又挂了一台欧洲机器。所以必须在 SparkSession 配置里显式固定时区而不是依赖集群环境。5.5 测试集 AUC 高上线后效果暴跌训练样本的幸存者偏差现象随机森林测试集 AUC 有 0.85上线后第二周风险分的均值就开始往上漂模型输出明显偏保守。原因训练数据用的是“已结案”的理赔案件但线上处理的很多是正在查勘中的未结案案件。已结案案件的理赔金额和责任认定都已经确定未结案案件的特征分布完全不同。这相当于训练集和线上分布不一样AUC 再高也是自欺欺人。解决取训练样本时只保留报案时间距今 90 天以上且已结案的案件这样给未结案案件留出了结案周期。同时要把“报案时间”作为特征而不是只看结案数据。上线后每天监控风险分的分布和正例占比一旦发现分数均值偏移超过设定阈值就触发重新训练。这一点看起来简单但在真实项目里是最容易被忽略的因为你看到的历史数据全都是“过去式”线上永远在等未来。6. 从模型上线到效果验证两个必看的评估指标和监控习惯6.1 用这种评估方式决定是否上线AUC 描述的是整体排序能力但保险公司真正关心的是“从一百万条案件里挑出最可疑的一万条能命中多少条”。这时候要算召回率在 top 千人甚至万人里的表现也就是precisionk。我一般会把模型输出的风险分降序排列取前 5% 作为高风险名单然后统计真实高风险样本结案后被确定为欺诈或高赔的案件在这份名单里的覆盖率。如果 top 5% 的覆盖率达到 60% 以上这个模型才值得推到生产环境。否则 AUC 再高业务侧用起来也没体感。from pyspark.sql.functions import col, percent_rank scored loaded_model.transform(feature_df) \ .withColumn(risk_rank, percent_rank().over(Window.orderBy(col(probability).desc()))) top_5 scored.filter(col(risk_rank) 0.05) hit_rate top_5.filter(col(is_high_risk) 1).count() / scored.filter(col(is_high_risk) 1).count() print(top 5% hit rate:, hit_rate)6.2 上线后每天盯一个指标模型上线后我给自己定的习惯是每天早上看一张报表当天案件的风险分均值、标准差和高分段占比。这三个数任何一个出现跳跃都要立刻查数据源是不是换了字段格式或者运营调整了报案流程。有一次我发现风险分均值在三天里涨了 20%查到最后是一条上游数据流把“报案时间”格式从北京时区变成了 UTC导致时间特征整体偏移。这种问题靠 AUC 验证是发现不了的只有实时监控能兜住。做完了清洗、建模、调优和监控你会明白保险行业的 Spark 项目最值钱的不是模型算法而是对数据分布变化的感知。一点经验每次重训模型前先跑一遍监控报告确认线上数据分布没有发生结构性变化再动手。希望帮到你。本文还有配套的精品资源点击获取
阅读完成 · 觉得有帮助?
咨询建站