
1. 项目背景与核心挑战数美科技作为一家专注于在线业务风控的技术公司每天需要处理数百TB级别的用户行为数据。这些数据主要来自各类互联网平台的实时交互日志、设备指纹信息和业务操作记录。在早期架构中我们面临三个典型痛点查询响应慢简单的用户行为分析需要1天以上的处理时间数据孤岛严重不同业务线的数据存储格式和访问方式不统一开发效率低下每个新需求都需要数据工程师重写ETL管道2. 技术架构演进路径2.1 原始架构分析最初的Lambda架构包含以下组件graph LR A[Kafka] -- B[Spark Streaming] A -- C[Flink] B -- D[HBase] C -- E[MySQL] D -- F[Presto] E -- F主要问题体现在实时/离线链路割裂多套存储引擎维护成本高跨系统JOIN性能差2.2 新一代平台设计我们采用流批一体统一存储的思路重构架构graph TB A[Kafka] -- B[Spark Structured Streaming] B -- C[Delta Lake] C -- D[ClickHouse] C -- E[Presto]关键改进点所有数据统一进入Delta Lake作为唯一可信源使用Spark SQL作为统一计算引擎按场景分流ClickHouse服务OLAPPresto支持ad-hoc查询3. 核心技术创新3.1 JSON数据处理优化原始日志90%以上是嵌套JSON我们开发了智能Schema推断系统// 自动识别JSON schema并注册为临时视图 val rawDF spark.read.format(json).load(path) val inferredSchema SchemaInferencer.infer(rawDF) spark.catalog.createTempView(logs, rawDF) // 动态生成查询SQL val query sSELECT ${inferredSchema.analysisFields.mkString(,)} FROM logs对比测试结果处理方式1GB数据耗时CPU利用率传统解析78s220%优化方案12s150%3.2 ClickHouse调优实践针对风控场景的典型查询模式大量点查少量聚合我们做了以下优化1. 表引擎选择CREATE TABLE risk_events ( dt Date, event_time DateTime, user_id String, -- 其他字段... INDEX uid_idx user_id TYPE bloom_filter GRANULARITY 3 ) ENGINE ReplacingMergeTree PARTITION BY dt ORDER BY (user_id, event_time)2. 冷热数据分层热数据SSD存储ReplicatedMergeTree温数据HDD存储普通MergeTree冷数据自动归档到对象存储4. 平台能力输出4.1 查询体验提升通过统一元数据管理和查询网关实现了标准SQL接口兼容自动路由到最优执行引擎查询结果缓存典型查询响应时间对比查询类型旧平台新平台用户行为轨迹6h30s设备关联分析2h45s跨业务统计失败2m4.2 开发模式变革引入Data as Code理念# 查询定义示例 source: risk_events metrics: - name: fraud_count expr: countIf(statusfraud) dimensions: - province - device_type filters: - dt 2023-01-01开发效率提升新需求交付周期从3天缩短到2小时业务人员自主完成80%的分析需求5. 踩坑经验5.1 Spark调优要点executor配置避免小文件问题spark.sql.shuffle.partitions2000 spark.executor.memoryOverhead2gDelta Lake优化OPTIMIZE delta./path/to/table ZORDER BY (user_id)5.2 ClickHouse注意事项避免高频小批量INSERT建议批次≥1000行合理设置max_threads建议物理核心数60%监控Merge操作频率及时调整分区策略6. 未来规划当前正在验证的方向基于GPU加速的实时图计算自适应压缩算法按列动态选择压缩格式智能预聚合自动识别高频查询模式关键认知大数据平台建设不是单纯的技术选型而是要构建符合业务节奏的数据供应链。我们的实践表明当查询延迟从天级降到分钟级时业务方会产生全新的分析需求进而推动数据应用的质变。