ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hadoop+Spark信贷风控系统:从架构设计到本地部署实战

Hadoop+Spark信贷风控系统:从架构设计到本地部署实战 简介面向金融信贷风控场景的大数据工程毕设资源基于 Hadoop 与 Spark 技术栈实现适合计算机、大数据、人工智能等专业的毕业设计选题、课程设计及项目实训。资源包含完整的系统设计与源代码整体规模为69个文件、约58KB其中以 Java 业务代码与 Scala 计算逻辑为主辅以 XML 配置文件、properties 参数配置、SQL 建表脚本和 README 说明文档可按模块快速理清数据处理、风险评分与信贷审批流程。目前已有325人学习浏览代码经过环境测试可稳定运行并支持下载后远程咨询与讲解。通过该资源可重点学习海量信贷数据批次处理与实时计算思路、Hadoop/Spark 集群任务调度方法、风险指标建模及后端接口设计是一份便于二次开发与毕业设计答辩展示的完整参考工程。1. 基于 Hadoop、Spark 的信贷风控系统先跑通再谈架构信贷审批如果靠人工逐笔看流水和征信一天处理不了几百单换成 Hadoop 存全量数据、Spark 跑特征和评分单子再多也能在分钟级出结果。这套基于 Hadoop、Spark 的大数据金融信贷风险控系统就是一套把离线批处理和实时流计算串起来的完整工程credit-risk-control 主业务模块负责进件和规则校验data-source-spark-streaming 负责实时接数h5-credit-risk-control 是申请端页面databases 目录里放着建表脚本。源码是 Maven 多模块结构IDEA 打开就能逐模块跑适合做毕设、课程设计也适合想搞清楚“真实风控数据链路到底怎么组织”的大数据开发。2. 先看架构再碰代码这条风控数据链路的四个关键节点2.1 为什么信贷风控要上 HadoopSpark而不是纯 MySQL信贷风控的特征数据远比一般业务系统复杂。用户基础信息、银行卡流水、多头借贷记录、电商消费行为、黑名单命中情况这些数据来源不同、更新频率不同而且有个共同点写一次、读无数次、越攒越久。用 MySQL 硬扛一是存储成本撑不住二是亿级表上的聚合 SQL 跑不动。Hadoop 的 HDFS 把半结构化日志和 CSV 全量囤下来Hive 负责离线宽表加工Spark 则擅长把几亿条记录的内存计算压缩到分钟级。这个项目也是这么分的用户画像、月度负债比、历史逾期率这类不苛求实时的指标走离线批处理进件瞬间的黑名单命中、设备指纹异常、短时间申请频次则走 Spark Streaming。两条线并行最终都落到业务库给审批接口用。这套分层放在真实数仓里是有明确血缘的。ODS 层直接存 Kafka 同步过来的进件日志和流水日志文件落在 HDFS 的 /user/hive/warehouse/ods_apply_log 这类路径下DWD 层做清洗和拉宽把用户、订单、还款、逾期事实拆成明细事实表DWS 层按天汇总用户粒度的负债比、查询次数、逾期率。后面评分用的特征基本都是从 DWS 层宽表里直接 select。看不懂代码时先问自己一句这条数据现在在哪个层定位能快一半。我为什么强调先看架构因为信贷风控系统的难点不在某个算法而在数据怎么按时、按序、按正确粒度汇到一起。如果一上来就盯某个类的实现很容易在局部打转。先花半小时把数据流捋顺后面改任何一块都知道会影响谁。反过来说遇到评分结果对不上也能按数据分层逐级排查而不是在业务代码里瞎找。2.2 源码包里的模块分别负责什么解压 zip 后第一眼会看到一堆文件和目录。很多人会懵“为什么有两个 pom.xml为什么有前端目录”其实这是 Maven 多模块工程。我建议别急着翻代码先把下面几个模块对号入座目录 / 模块技术角色实际职责credit-risk-controlSpring Boot 后端进件接口、规则校验、评分结果查询>{ userId: U10002345, deviceId: D8671E04A, applyTime: 1733827200000, applyId: A202401010001, loanAmount: 50000, term: 12 }第三步data-source-spark-streaming 消费该 topic在内存里做滑窗统计同一身份证 1 小时内申请次数、同一设备号 10 分钟内申请次数、同一手机号当天关联申请数。这些是典型的反欺诈实时指标。第四步离线批处理任务定期从 HDFS 读历史借贷流水计算负债收入比、近 6 个月逾期次数、额度使用率等强变量连同实时特征一起组装成评分特征向量。第五步规则引擎或模型给出风险评分与拒绝/人工审核/通过建议结果回写 risk_result 表H5 端轮询查询到最终状态。这套结构最巧妙的地方是实时和离线解耦。Spark Streaming 只算短窗口里的频次特征不需要查全量历史离线任务专啃大表。两者互不抢资源又能通过 apply_id 或 user_id join 到一起。调试时也可以分开验证实时链路只看 risk_mark 字段离线链路只看 credit_score 字段谁出了问题都不用把整个系统停掉。2.4 伪分布式环境下的最小部署矩阵架构听明白了接下来要解决“我这台电脑能不能跑”。我拆过不少大数据毕设最怕的就是用户一上来搭三台虚拟机结果内存爆掉。8G 内存的笔记本完全够用关键是 Hadoop 用伪分布式、Spark 用 local[*]、Kafka 单节点、MySQL 本地。下面是我用的最小部署矩阵组件部署方式建议内存关键配置Hadoop HDFS伪分布式1.5Gdfs.replication1Hadoop YARN同节点1Gyarn.nodemanager.resource.memory-mb2048Sparklocal[*] 或 standalone2Gspark.executor.memory1gKafka单节点512Mlog.retention.hours24MySQL本地512M默认端口 3306Zookeeper单节点256M默认端口 2181内存分配逻辑很简单HDFS 的 DataNode 和 NameNode 各占一部分YARN 给 2GSpark 如果也跑在同一台机器上就不要再独占太多。伪分布式模式下 HDFS 的 core-site.xml 可以这样配configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/usr/local/hadoop/tmp/value /property /configurationfs.defaultFS 指明了 NameNode 地址hadoop.tmp.dir 必须设成可写目录否则 namenode format 时会报权限错误。YARN 的 yarn-site.xml 里我一般会关掉资源强校验避免容器申请内存大于物理内存时直接把任务 kill 掉property nameyarn.nodemanager.vmem-check-enabled/name valuefalse/value /property property nameyarn.nodemanager.resource.memory-mb/name value2048/value /propertyvmem-check-enabled 是伪分布式环境最容易忽略的一项。默认开启时Spark 任务经常报 Container killed by YARN for exceeding memory limits关闭后就能跑。这算是我踩过的血泪经验。HDFS 启动后还要记得建好 Hive 数仓需要的目录不然后面离线任务往不存在路径写数据会直接抛异常hdfs dfs -mkdir -p /user/hive/warehouse/ods_apply_log hdfs dfs -mkdir -p /user/hive/warehouse/dws_user_feature hdfs dfs -chmod -R 777 /user/hive路径名本身就是分层语义ods 放原始日志dws 放用户汇总特征。目录权限给 777 在单机测试环境没问题生产环境不要这么干那是另一个安全话题。2.5 实时与离线产出的特征怎么合并信贷风控的评分不能只靠实时频次也不能只靠月度离线画像。合并方式通常是Spark Streaming 算出的实时特征写入 Rediskey 设为 apply_id 或 user_idTTL 设置为 24 小时离线任务每天凌晨跑出宽表存到 Hive再同步一份到 MySQL/ES在线评分时后端同时查 Redis 和 MySQL把两条特征拼接成完整向量。这套源码里也保留了类似设计实时模块输出 risk_mark离线模块输出 credit_score最终在 risk_result 表里按进件申请号合到一起。合并时最容易出的问题是对不上时间口径。实时特征是“截至当前时刻的前 1 小时”离线特征是“截至昨天 24 点的完整月份”两者本来就不在同一时间平面写评分逻辑时不要试图让它们严格相等而是让离线数据作为基准实时数据做增量修正。举个例子用户昨天负债比是 0.4今天又申请了一笔消费贷实时模块只把“近 1 小时申请次数加 1”追加进去而不能把离线特征里的总负债直接改掉。时间口径对不上评分结果就是错的这在信贷业务里不是小事。3. 本地跑通这套信贷风控系统环境、顺序与启动参数3.1 环境准备先装什么后装什么版本怎么定拿到源码后第一件事不是导入 IDEA而是把环境装到“能跑”的状态。这套项目依赖 JDK、Maven、Hadoop、Zookeeper、Kafka、Spark、MySQL 和 Node。按依赖关系我的安装顺序是JDK 8 → MySQL 5.7/8.0 → Zookeeper → Hadoop → Spark → Kafka → Node 14。JDK 版本别乱升很多 Spark 2.x 的 jar 在 JDK 11 下会报模块访问错误如果项目 pom 里锁定的是 Spark 2.4建议用 JDK 8 一条路走到底。Hadoop 和 Spark 的版本匹配是个玄学点。Hadoop 2.7/2.8 和 Spark 2.3/2.4 是经典组合Hadoop 3.x 则更适合 Spark 3.x。源码的 pom.xml 里如果已经写了版本号就以它为准如果没写我会用 Hadoop 2.7.7 Spark 2.4.8这个组合的文档最多遇到问题也最容易搜到答案。Zookeeper 用 3.4.14Kafka 用 2.11 对应版本。注意 Kafka 2.11 这里的 2.11 是 Scala 编译版本不是 Kafka 版本号这个坑后面还会提。如果你不想在本机装全套也可以考虑直接用 Docker 镜像跑伪分布式 Hadoop但我个人建议第一次调试还是本机直接装因为能看到完整日志。Docker 虽然快日志和端口映射多了一层遇到问题排查成本反而高。环境变量方面至少要确保这些变量都生效export JAVA_HOME/usr/local/jdk1.8 export HADOOP_HOME/usr/local/hadoop export SPARK_HOME/usr/local/spark export PATH$PATH:$HADOOP_HOME/bin:$SPARK_HOME/bin如果是在 Windows 上调试还要额外加一个 HADOOP_HOME 指向包含 winutils.exe 的目录这个问题在避坑章节详细说。3.2 数据库初始化先建库再谈启动后端接口和流任务都要读业务表所以数据库必须最先初始化。打开 databases 目录里面应该有建表脚本。如果你用的是 MySQL 客户端直接执行mysql -uroot -p123456 databases/credit_risk.sql这里的 -p 后面跟的是本地 MySQL root 密码如果你本机密码不是 123456就改成自己的。执行成功后可以验证一下SHOW TABLES; SELECT COUNT(*) FROM user_info;我一般会先看 user_info 和 apply_record 这两张表有没有初始数据。很多毕设项目的 SQL 脚本里只建表不插数后续接口调试时查不到数据容易被误判成程序 bug。如果发现脚本只建表没造数自己补一条测试用户即可别去改业务代码INSERT INTO user_info (user_id, user_name, id_card, mobile, monthly_income, create_time) VALUES (U10002345, 测试用户, 110101199001011234, 13800138000, 12000, NOW());执行 SQL 时还有一个常见问题字符集。如果表结构里定义了 utf8mb4但 MySQL 服务端默认字符集是 latin1执行中文注释会直接报错。启动 MySQL 时加上 --character-set-serverutf8mb4 或在 my.cnf 里配置能少很多事。连接串里的 serverTimezone 也要配成 Asia/Shanghai否则 Spring Boot 起来后查时间字段会差 8 小时。3.3 导入 IDEAMaven 多模块要按这种方式打开这个项目不能用“File → Open”随便选一层目录因为两个业务模块共用一套父 pom。正确做法是在 IDEA 里 Open 到含根 pom.xml即 CreditRiskControl.iml 所在目录那一层让 IDEA 识别为 Maven 多模块工程。导入时注意 JDK 选 8Maven 的 settings.xml 里配阿里云镜像否则 spark-streaming-kafka 相关依赖可能下载很慢。如果你是命令行党也可以不进 IDEA直接用 Maven 构建mvn clean package -DskipTests-DskipTests 表示跳过单元测试只打 jar。第一次跑命令会下载大量依赖看到 BUILD SUCCESS 才算完。如果报依赖下载失败先检查 settings.xml 的 mirror 是否生效再检查本地仓库 .m2 目录是不是被什么工具锁了权限。我还会习惯性地在导入后执行一次mvn dependency:tree这个命令会列出所有传递依赖。它最大的价值不是看有哪些包而是排查冲突。如果在树里看到同一个 jar 出现两个版本号就要注意避坑章节里提到的 NoSuchMethodError 了。多模块导入完成后IDEA 右侧 Maven 面板里应该能看到 credit-risk-control 和>hadoop namenode -format start-dfs.sh jpsjps 输出里要能看到 NameNode 和 DataNode。如果你用的是伪分布式第一次启动前必须 format否则 NameNode 会一直报 Incompatible namespaceIDs。第二步启动 YARNstart-yarn.sh jps这次应该多出 ResourceManager 和 NodeManager。第三步启动 Spark 历史服务方便查看任务执行细节$SPARK_HOME/sbin/start-history-server.sh第四步验证 Spark 和 Hadoop 的集成spark-shell --master yarn --deploy-mode client能进 spark-shell 说明 Spark 能正常向 YARN 申请资源。之后再启动 Zookeeper 和 Kafka创建业务所需的 topickafka-topics.sh --create --zookeeper localhost:2181 --topic apply-events --partitions 3 --replication-factor 1topic 的分区数会影响 Spark Streaming 的并行度。单机伪分布式下分区数设为 1 或 3 都可以如果设成 3Spark 端对应的 repartition 也要是 3否则会出现无用 shuffle。验证 Kafka 是否正常可以用生产消费一条消息而不是只看进程列表kafka-console-producer.sh --broker-list localhost:9092 --topic apply-events kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic apply-events --from-beginning看到能收发消息实时链路才算通。很多人在这一步直接跳过结果后面 Spark 里查不到数据还得回头来查 Kafka。3.5 Spark Streaming 任务的两种启动姿势流任务不是 Spring Boot 那种常驻服务它的启动方式有讲究。本地调试时我直接在 IDEA 里运行>spark-submit \ --class com.credit.risk.streaming.StreamingRiskJob \ --master yarn \ --deploy-mode client \ --executor-memory 1g \ --num-executors 1 \ >java -jar credit-risk-control.jar --spring.profiles.activedev前端 H5 在 h5-credit-risk-control 目录下先装依赖再启动npm install npm run devnpm run dev 默认监听 8080 之类的端口。如果和后端端口冲突改 package.json 里的 dev 脚本参数或者把后端 server.port 改掉别让两个进程挤在同一个端口上这是新手最容易漏的一步。4. 风控计算核心实时频次、离线画像与评分规则怎么落地4.1 Kafka 消费端从 JSON 报文到安全的风险事件实时模块的第一步是消费 Kafka 里的申请事件。下面这段代码是典型的 Spark Streaming Kafka 消费端写法项目里>val kafkaParams Map[String, Object]( bootstrap.servers - localhost:9092, key.deserializer - classOf[StringDeserializer], value.deserializer - classOf[StringDeserializer], group.id - credit-risk-group, auto.offset.reset - latest, enable.auto.commit - (false: java.lang.Boolean) ) val stream KafkaUtils.createDirectStream[String, String]( streamingContext, PreferConsistent, Subscribe[String, String](Array(apply-events), kafkaParams) ) val riskEvents stream.map(record { val json parse(record.value()) val userId (json \ userId).extract[String] val deviceId (json \ deviceId).extract[String] val applyTime (json \ applyTime).extract[Long] (userId, deviceId, applyTime) })这段代码里有三个关键点。第一auto.offset.reset 设置为 latest表示新任务启动时只消费启动之后的增量消息这在联调时很实用如果你想回放历史消息改成 earliest。第二enable.auto.commit 必须设为 false因为我们要等业务处理完再手动提交 offset防止任务在写结果表前崩溃导致数据丢失。第三PreferConsistent 指定了分区分配策略它会让每个 executor 尽量均匀拿到分区避免某个节点过载。实际项目中可能还会在 map 之前加 filter把主题消息里不是申请事件的日志、心跳数据过滤掉避免下游解析 JSON 报错。反序列化这一层看着简单其实是整个流处理的咽喉。我用过一个血泪经验后端接口改了字段名但流任务里的 JSON 解析没有同步改导致评分特征全是默认值而且不报错。所以做强类型解析时尽量给每条消息加一个 schema 版本号解析逻辑向前兼容。4.2 滑窗统计同一身份证一小时内申请几次算风险实时反欺诈最常用的手段就是滑窗计数。统计短时间内的申请次数超过阈值就标记为高风险。核心代码大致长这样riskEvents .map(event (event._1, 1L)) .reduceByKeyAndWindow( (a: Long, b: Long) a b, (a: Long, b: Long) a - b, Seconds(3600), Seconds(60) ) .filter(_._2 5) .map { case (userId, count) RiskMark(userId, HIGH_FREQ_APPLY, count) }reduceByKeyAndWindow 的前两个参数分别是窗口内加法函数和窗口滑出时的减法函数。Seconds(3600) 是窗口长度表示统计“过去 1 小时”Seconds(60) 是滑动步长表示每 60 秒计算一次。这里用减法函数会比每次全量重算快很多Spark 会保存中间状态这也是窗口计算能跑得动的原因。阈值 5 是测试时设的初值真实环境要根据业务容忍度调现金贷产品可能 3 次就触发大额抵押贷可以放宽到 10 次。阈值不建议写死在代码里后面的进阶章节会说明怎么改成动态配置。滑窗计算的常见翻车点是数据倾斜。如果某几个 user_id 的申请频次特别高导致这些 key 集中在同一个 executor会造成 OOM。我在生产环境见过一个中介机构用同一批手机号批量进件把某个 executor 的内存直接打爆。解决办法是在 key 上加盐做预聚合或在源头限制单用户并发两种方案都能缓解但别指望 Spark 自动帮你分好区。计算结果除了落库还应该写一份到 Redis方便后端查询实时风险标记redisClient.setex(srisk:realtime:${userId}, 86400, riskMark.toString)86400 是 TTL单位秒表示这个实时标记只在当天有效。信贷审批讲究时效性昨天的实时频次对今天的进件已经没有意义所以 TTL 一定要设置而不是让 key 永久留在 Redis 里。4.3 离线特征与贷前评分Spark SQL 或 Hive SQL 做宽表实时频次只能拦住“短时间重复申请”这类明显的欺诈。真正的信用风险评估要靠离线特征收入负债比、历史逾期率、信用卡使用率、贷款查询次数。这些指标都要用历史流水算。离线任务读取 HDFS 上的流水表和还款表用 Spark SQL 做特征加工。典型的 SQL 思路如下SELECT u.user_id, round(sum(l.loan_amount) / nullif(u.monthly_income, 0), 4) AS debt_income_ratio, count(if(r.status OVERDUE, 1, NULL)) / nullif(count(r.repayment_id), 0) AS overdue_rate, round(avg(c.card_balance / nullif(c.card_limit, 0)), 4) AS credit_utilization FROM user_info u LEFT JOIN loan_apply l ON u.user_id l.user_id LEFT JOIN repayment_record r ON u.user_id r.user_id LEFT JOIN credit_card c ON u.user_id c.user_id WHERE l.apply_time date_sub(current_date(), 180) GROUP BY u.user_id这段 SQL 里有三个常用的特征变量。debt_income_ratio 是负债收入比分子是近 6 个月累计放款金额分母是月收入nullif 用来防止除数为零。overdue_rate 是逾期率它把逾期记录数除以总还款笔数这个变量在风控模型里权重通常很高。credit_utilization 是信用卡额度使用率也是经典的风险因子。注意 LEFT JOIN 的顺序如果业务上某些用户没有任何贷款或信用卡记录LEFT JOIN 能保证用户不丢只是特征值为空后续再用均值或 0 填充。这段 SQL 在 Spark SQL 和 Hive 里都能跑但要注意 date_sub 函数的兼容性。Spark 2.4 的 date_sub 返回 DateType直接跟字符串比较有时会隐式转换失败稳妥做法是先 cast 成 date 再比较WHERE l.apply_time cast(date_sub(current_date(), 180) as date)离线任务的触发方式毕设环境用 Linux crontab 就可以满足0 2 * * * /usr/local/spark/bin/spark-submit --class com.credit.risk.batch.OfflineFeatureJob credit-risk-control-1.0-SNAPSHOT.jar0 2 * * * 表示每天凌晨两点跑一次。凌晨跑批的好处是业务低峰期资源和数据库压力都小。如果将来任务多了再换 Azkaban 或 DolphinScheduler 做工作流编排这个项目阶段不需要。4.4 评分输出结果表怎么设计才能方便审批端查询实时和离线特征算完后评分结果要落库。risk_result 表的设计直接影响查询性能我在项目里见过把十几列特征全塞进一张宽表的做法结果每次查询都要扫一大片。更合理的做法是只保留最终评分、风险等级、命中规则和特征版本字段类型说明apply_idvarchar(64)进件编号主键user_idvarchar(32)用户编号credit_scoreint最终评分范围 0-100risk_levelvarchar(16)LOW / MEDIUM / HIGHhit_rulesvarchar(512)命中的规则列表逗号分隔score_versionvarchar(16)评分模型版本create_timedatetime计算时间后端审批接口只需要按 apply_id 查这一条记录就能决定自动通过、人工审核还是拒贷。如果你要做模型迭代score_version 字段能帮你回溯某段时间的结果是哪个模型产的这个字段在初版设计时最容易漏。查询时配合索引效率会好很多ALTER TABLE risk_result ADD INDEX idx_user_id (user_id);别小看这个索引。如果审批端要查“某个用户的所有申请记录”没有索引的话随着表数据增长查询会越来越慢。大数据系统里最后一道查询路径往往是最容易被忽视的性能瓶颈。5. 避坑指南HadoopSpark 跑批常见的五个故障下面这些坑是我在本地跑通这类毕设和真实项目时都遇到过的按出现频率排序。每一条都按“现象 → 原因 → 解决”写清楚。5.1 启动 Spark 任务就报 NoSuchMethodError现象spark-submit 提交任务后几秒内就抛 NoSuchMethodError指向某个 netty 或 protobuf 类。原因Hadoop、Spark 依赖的第三方库版本不一致。比如 Spark 2.4 内置 netty 3.x但 pom.xml 里引了 netty 4.xclasspath 里靠前的那个类把另一个覆盖了。解决把用户代码 pom 里跟 Hadoop/Spark 冲突的依赖全部标成 provided让它们运行时用集群自带的版本dependency groupIdorg.apache.spark/groupId artifactIdspark-streaming-kafka-0-10_2.11/artifactId version2.4.8/version scopeprovided/scope /dependencyscope 改为 provided 后IDEA 本地跑时仍会带包但 spark-submit 不会把这份 jar 打进最终包有效避免版本冲突。我每次排查这类报错第一件事是执行 mvn dependency:tree 看冲突路径而不是去改 Spark 源码。5.2 本地模式运行报找不到 winutils.exe现象Windows 上跑 Spark 或访问 HDFS报 Failed to locate the winutils binary in the Hadoop binaries。原因Spark 访问 HDFS 时通过 hadoop-common 里的 native 方法检测 Windows 环境而 Windows 没有对应的 hadoop.dll/winutils.exe。这不是代码问题是 Spark 在 Windows 下的环境兼容问题。解决下载对应 Hadoop 版本的 winutils.exe 放到某个目录并在环境变量里指定export HADOOP_HOMED:/hadoop-winutils然后把 hadoop.dll 所在的 bin 目录加到 PATH。这只影响本地联调部署到 Linux 服务器后就不存在。很多新手被这个报错吓到其实只是缺一个二进制文件。5.3 Kafka 里能看到消息Spark Streaming 却消费不到现象kafka-console-consumer 能正常收到消息但 Spark Streaming 端一直不打印业务日志。原因最常见的是消费组 offset 已经提交到很后面的位置而 Spark 任务配置的 auto.offset.reset 是 latest从当前最新位置开始消费自然错过之前发的旧消息。解决先把 auto.offset.reset 改成 earliest或者用 kafka-consumer-groups.sh 把消费组 offset 重置到开头再启动任务。另一个隐蔽原因是入参顺序不对代码里读的参数顺序是 broker 地址、topic 名、group id而你 spark-submit 传参时写反了 topic 和 group id导致任务订阅了错误的 topic。这种问题不会报错只能靠日志确认实际订阅的是哪个 topic。检查方法是在代码里打印一句日志INFO: Subscribe topic apply-events, group credit-risk-group如果日志显示的不是预期值回头对一下传参顺序这是最容易忽略的细节。5.4 批次处理时 OOM而集群内存明明够现象处理大批量数据时 executor 持续 GC频繁 Full GC 后直接 OOM。原因Spark Streaming 默认批次大小和限速参数没调。生产环境里 Kafka topic 分区数大于 executor 数量时每个 executor 要同时拉取多个分区的数据积压在内存里另外流任务默认没有开启背压消费速度不受控制。解决开启背压并设置合理的速率streamingContext.conf.set(spark.streaming.backpressure.enabled, true) streamingContext.conf.set(spark.streaming.backpressure.initialRate, 1000) streamingContext.conf.set(spark.streaming.kafka.maxRatePerPartition, 500)backpressure.enabled 让 Spark 根据处理速度动态调节消费速率initialRate 是初始每分区每秒最大消费条数maxRatePerPartition 是硬上限。没有这些参数任何 Spark Streaming 应用在高峰期都有 OOM 风险。这也解释了为什么同一套代码在演示环境很顺畅、一到真实流量就挂。5.5 前后端页面看到申请记录但迟迟不出评分现象H5 页面往后端查申请状态能看到申请记录但 risk_result 一直为空。原因这条记录的实时风险标记没有写回或写回了但 apply_id 对不上。最常见的情况是实时任务启动时日志里能看到数据流但结果写库环节没有正确处理空值另一个可能是实时任务和离线任务写的是同一张结果表后启动的任务在 create_time 上加了自己的默认值把另一条记录覆盖了。解决先用 SQL 查这条进件的原始申请事件在 Kafka 里是否真的存在确认存在后再检查实时任务里写入 result 表的 SQL 是不是用了 insert or ignore。如果表里有 apply_id 唯一索引重复写入静默失败不会报错这个问题会花不少时间才能定位。我通常会在实时任务的结果写入处加一条日志输出写入影响行数。别小看这一行日志它比什么排查工具都管用尤其是任务看起来一切正常但数据就是不对的时候整个过程就像在查一个黑匣子。6. 进阶用法把评分规则从硬编码改成可动态调整的配置系统跑通后下一步不是去加新功能而是把写死的阈值变成配置。因为信贷风控规则是天天变的上线早期审批可以松逾期上来了就要收紧。如果每个阈值都改代码重新打包速度太慢而且容易把不同环境的 jar 搞混。我习惯的做法是把规则存进 MySQL 的一张表实时任务每次启动时加载同时每隔几分钟刷新一次。规则表结构很简单字段类型说明rule_codevarchar(32)规则编码rule_paramvarchar(32)参数名rule_valuevarchar(128)参数值effective_timedatetime生效时间expire_timedatetime失效时间比如之前代码里写死的 1 小时 5 次申请可以拆成两条配置INSERT INTO rule_config VALUES (HIGH_FREQ_APPLY, window_seconds, 3600, 2024-01-01 00:00:00, 2025-12-31 23:59:59), (HIGH_FREQ_APPLY, max_count, 5, 2024-01-01 00:00:00, 2025-12-31 23:59:59);然后把 4.2 节里那段滑窗代码改成从配置读取参数val windowSeconds getRuleConfig(HIGH_FREQ_APPLY, window_seconds).toInt val maxCount getRuleConfig(HIGH_FREQ_APPLY, max_count).toInt riskEvents .map(event (event._1, 1L)) .reduceByKeyAndWindow( (a: Long, b: Long) a b, (a: Long, b: Long) a - b, Seconds(windowSeconds), Seconds(60) ) .filter(_._2 maxCount)这一改动看起来只多了两行价值却很大规则调整不再需要重新打包和重启流任务只要在管理页面改数据库最多等配置刷新周期结束就能生效。我用这个方案支撑过多个产品线的差异化风控策略。实际改造时要注意两点。第一别把规则加载写在 map 算子里面否则每个分区每批次都要查库一次数据库压力会很大。正确做法是在 Driver 端定时加载配置然后通过广播变量把配置传给 executor。广播变量在 Spark Streaming 里尤其有效因为 executor 数量少、配置变更频率低。第二改配置不是立刻生效要留一个生效延迟的预期并记录每次规则的调整历史。信贷风控有审计要求哪天一门心思调阈值却忘了留痕后面出问题很难追溯。从那以后我每次接手这类风控项目都会先问一句“规则在哪配”如果答案不是“配置表”我会建议业务和技术一起出个方案尽早把硬编码拆掉。规则引擎的价值不在复杂度而在能让你在业务变化时不动代码。希望帮到你。本文还有配套的精品资源点击获取
返回列表