news 2026/8/10 3:18:36

数美科技大数据平台:从Hadoop到ClickHouse的实时风控演进

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
数美科技大数据平台:从Hadoop到ClickHouse的实时风控演进

1. 数美科技大数据平台演进概述

数美科技作为国内领先的在线业务风控服务商,其大数据平台经历了从传统批处理到实时交互式分析的完整演进过程。早期平台采用典型的Hadoop生态架构,数据查询响应时间普遍超过24小时,严重制约了业务决策效率。而当前平台已实现"定义即可查"的实时交互能力,查询延迟降低至秒级,数据规模突破数百TB级别。

这个转变背后是三个核心突破点:首先是存储引擎从HDFS迁移到ClickHouse列式数据库,使压缩比提升5倍的同时查询性能提高10倍以上;其次是计算框架引入Spark SQL实现批流统一处理,告别了原先MapReduce的繁重开发模式;最后是采用JSON作为统一数据交换格式,简化了上下游系统集成复杂度。

2. 技术架构深度解析

2.1 存储层:ClickHouse的极致优化

在ClickHouse集群部署上,数美采用双副本分片集群模式,每个分片由6节点组成。针对风控场景的特殊需求,对默认配置进行了多项关键调整:

-- 建表示例展示特殊配置 CREATE TABLE risk_events ( event_time DateTime64(3, 'Asia/Shanghai'), user_id String CODEC(ZSTD(3)), device_fp String CODEC(ZSTD(5)), risk_score Float32, features Nested( name String, value Float32 ) ) ENGINE = ReplicatedMergeTree() PARTITION BY toYYYYMM(event_time) ORDER BY (user_id, event_time) TTL event_time + INTERVAL 180 DAY SETTINGS index_granularity = 8192;

关键配置说明:ZSTD压缩算法针对不同字段设置差异化的压缩级别,设备指纹等长文本采用更高压缩比;TTL机制自动清理过期数据;DateTime64精确到毫秒满足风控时序需求。

实际运行中遇到的最大挑战是稀疏列存储问题。当某些特征字段空值率超过90%时,默认存储方式会浪费大量空间。通过启用allow_suspicious_low_cardinality_types参数,对枚举型字段使用LowCardinality类型,使存储空间减少40%。

2.2 计算层:Spark SQL的批流统一

数美在Spark集群部署上选择YARN资源管理模式,采用动态资源分配策略:

spark-submit \ --master yarn \ --deploy-mode cluster \ --conf spark.dynamicAllocation.enabled=true \ --conf spark.shuffle.service.enabled=true \ --conf spark.dynamicAllocation.maxExecutors=100 \ --conf spark.sql.adaptive.enabled=true \ --executor-memory 16G \ --executor-cores 4 \ --class com.shumei.risk.AnalysisJob \ risk-analysis.jar

特征工程中大量使用Spark SQL的窗口函数实现滑动时间窗统计:

val behaviorFeatures = spark.sql(""" SELECT user_id, COUNT(*) OVER ( PARTITION BY user_id ORDER BY event_time RANGE BETWEEN INTERVAL 1 HOUR PRECEDING AND CURRENT ROW ) AS hour_actions, AVG(risk_score) OVER ( PARTITION BY device_fp ORDER BY event_time ROWS BETWEEN 100 PRECEDING AND CURRENT ROW ) AS device_avg_score FROM risk_events WHERE event_time >= NOW() - INTERVAL 7 DAY """)

2.3 数据交换:JSON Schema治理

采用JSON Schema规范数据格式,使用ajv工具进行校验:

{ "$schema": "http://json-schema.org/draft-07/schema#", "type": "object", "properties": { "event_time": { "type": "string", "format": "date-time" }, "risk_score": { "type": "number", "minimum": 0, "maximum": 1 }, "ip_info": { "type": "object", "properties": { "country": {"type": "string"}, "asn": {"type": "integer"} } } }, "required": ["event_time", "risk_score"] }

在Spark中通过自定义SerDe实现JSON高效解析:

val schema = new ArrayType( StructType(Seq( StructField("feature_name", StringType), StructField("feature_value", FloatType) )) ) spark.read.schema(schema) .option("mode", "FAILFAST") .json("/data/events/*.json")

3. 性能优化实战

3.1 ClickHouse查询加速技巧

针对风控场景的典型查询模式,我们设计了特殊的物化视图:

CREATE MATERIALIZED VIEW risk_stats_hourly ENGINE = AggregatingMergeTree() PARTITION BY toYYYYMMDD(event_time) ORDER BY (user_id, feature_type) POPULATE AS SELECT user_id, featureType AS feature_type, toStartOfHour(event_time) AS hour_time, countState() AS event_count, avgState(risk_score) AS avg_score FROM risk_events GROUP BY user_id, featureType, hour_time;

查询时使用最终模式(FINAL)保证数据一致性:

SELECT user_id, sumMerge(event_count) AS total_actions, avgMerge(avg_score) AS overall_risk FROM risk_stats_hourly FINAL WHERE hour_time >= NOW() - INTERVAL 24 HOUR GROUP BY user_id HAVING total_actions > 10 ORDER BY overall_risk DESC LIMIT 1000;

3.2 Spark资源调优经验

通过分析历史任务得出最佳资源配置比例:

任务类型Executor数量单Executor内存并行度系数
特征计算5016G0.8
模型训练3032G0.5
数据清洗808G1.2

关键Spark参数设置经验:

  • spark.sql.shuffle.partitions设为集群核心数×并行度系数
  • 对于JOIN操作设置spark.sql.autoBroadcastJoinThreshold=64MB
  • 启用spark.speculation=true应对慢节点问题

4. 平台运维关键指标

4.1 集群监控体系

ClickHouse集群核心监控项:

指标名称预警阈值采集频率
ReplicatedQueueSize>100010s
MemoryUsage>90%5s
MergeSpeed<10MB/s1m
QueryDuration99分位>5s实时

使用Prometheus+Granafa构建的监控看板包含:

  • 查询QPS热力图
  • 内存使用趋势
  • 副本同步延迟
  • 慢查询TOP10

4.2 数据质量检查

每日例行数据质量检查项:

SELECT toDate(event_time) AS day, countIf(JSONIsValid(raw_data)=0) AS invalid_json, countIf(user_id='') AS empty_user, count() AS total, invalid_json/total AS bad_rate FROM risk_events WHERE event_time >= TODAY() - 1 GROUP BY day HAVING bad_rate > 0.0001;

5. 典型问题排查实录

5.1 ClickHouse内存溢出

现象:查询报错"Memory limit exceeded"但实际数据量不大 根因:GROUP BY字段基数过高导致中间状态爆炸 解决方案:

  1. 增加max_memory_usage参数
  2. 对高基数字段先做采样或分桶
  3. 使用approx_percentile替代精确计算

5.2 Spark数据倾斜

特征计算任务中某个stage耗时异常:

  1. 通过Spark UI定位到处理device_fp的task处理数据量是其他task的100倍
  2. 解决方案:
val skewedKeys = Seq("device123","device456") val broadcastKeys = spark.sparkContext.broadcast(skewedKeys) df.withColumn("is_skewed", when(col("device_fp").isin(broadcastKeys.value), true).otherwise(false)) .repartition(100, col("is_skewed"), rand())

6. 平台演进路线

当前架构仍面临两个主要挑战:首先是实时特征计算延迟需要从当前的5分钟降低到30秒以内,计划引入Flink替换部分Spark Streaming作业;其次是跨数据中心查询性能问题,正在测试ClickHouse的分布式表引擎。下一步重点将围绕以下方向:

  • 基于GPU加速的特征计算
  • 自适应查询优化
  • 自动化数据分级存储
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/10 3:18:25

护网行动(HVV)核心技术解析与实战指南

1. 护网行动&#xff08;HVV&#xff09;的本质与行业背景护网行动&#xff08;简称HVV&#xff09;是国内网络安全领域一项具有实战性质的攻防演练活动&#xff0c;最早可追溯至2016年由相关部门牵头组织。这项年度性安全演练最初主要覆盖重点行业单位&#xff0c;如今已发展成…

作者头像 李华
网站建设 2026/8/10 3:16:21

代码规范的价值与实践:提升团队协作与代码质量

1. 为什么我们需要代码规范刚入行那会儿&#xff0c;我最烦的就是看别人的代码。变量名全是a、b、c&#xff0c;缩进乱七八糟&#xff0c;有的地方用tab有的地方用空格&#xff0c;一个函数动辄几百行...每次接手这样的代码&#xff0c;我都想重写一遍。直到后来自己带团队&…

作者头像 李华
网站建设 2026/8/10 3:13:14

3D可视化技术在物流拼箱中的革命性应用

1. 项目背景&#xff1a;当传统拼箱遇上3D可视化革命 在物流和仓储领域&#xff0c;"多内盒混装拼箱"一直是个让人头疼的技术活。简单来说&#xff0c;就是把不同尺寸、不同形状的商品盒子&#xff0c;合理地装进一个标准尺寸的大箱子里。听起来容易&#xff1f;实际…

作者头像 李华
网站建设 2026/8/10 3:12:49

从工具调用到认知伙伴:Agent记忆系统的架构演进与实践路径

1. 从“工具调用者”到“认知伙伴”&#xff1a;Agent进化的分水岭最近和几个做AI应用的朋友聊天&#xff0c;发现一个挺有意思的现象&#xff1a;大家一提到“Agent”&#xff0c;脑子里蹦出来的第一反应&#xff0c;十有八九是“那个能调用API、执行任务的东西”。确实&#…

作者头像 李华
网站建设 2026/8/10 3:11:53

SpringBoot+Vue大学生迎新系统架构与优化实践

1. 项目概述&#xff1a;大学生迎新系统的技术架构与价值这套2025年最新版的大学生迎新系统采用SpringBootVue的前后端分离架构&#xff0c;搭配MyBatis持久层框架和MySQL数据库&#xff0c;是当前高校信息化建设中典型的轻量级解决方案。我在实际部署过三所高校的类似系统后发…

作者头像 李华