目录
一、引言:大数据与 AI 时代
随着互联网、物联网、5G 的迅猛发展,数据已成为继土地、劳动力、资本、技术之后的第五大生产要素。人类社会每天产生的数据量已从 GB、TB 时代迈入 PB、EB 时代。如何高效地存储、处理、分析这些海量数据,并从中提炼出价值,催生了大数据技术体系。而人工智能(尤其是以深度学习为代表的新一代 AI)的复兴,恰恰依赖于大数据提供的"燃料"——海量训练样本。大数据与 AI 的融合,构成了当代系统架构师必须掌握的核心技术栈。
· 大数据是"石油":提供海量、多源、高速的数据原材料
· AI 是"炼油厂":通过算法从数据中提炼模型与决策能力
· 二者相互成就:大数据让 AI 模型更准确,AI 让大数据价值更显著
· 架构师视角:需要同时理解数据管道(Data Pipeline)与模型管道(Model Pipeline),并将二者打通
1.1 为什么架构师必须理解大数据与 AI
在软考系统架构设计师的考察体系中,大数据与 AI 已经从"加分项"变为"必考点"。一个合格的架构师需要回答以下问题:
- 数据规模:当数据量从 GB 跃升到 PB 时,传统单机数据库架构为何失效?需要引入哪些分布式组件?
- 实时性:业务要求秒级响应,批处理架构无法满足,如何设计流式处理链路?
- 成本控制:海量存储与计算如何通过弹性伸缩、列式存储、数据分层降低 TCO?
- 模型迭代:模型从训练到上线、再到监控、再训练的闭环如何自动化?这正是 MLOps 要解决的问题。
- 合规与隐私:数据脱敏、联邦学习、隐私计算在架构中如何落地?
本文将围绕数据采集 → 存储 → 计算 → 分析 → 服务 → 模型训练 → 模型推理 → 模型治理这一完整链路,系统梳理大数据与 AI 架构的核心知识,并辅以大量 SVG 架构图、对比表和软考真题,帮助读者建立体系化的认知。
二、大数据基础
2.1 大数据的 5V 特征
大数据的定义性特征通常被概括为5V 模型,这是软考中最基础也最高频的考点之一,必须熟记每个 V 的中文含义与典型表现。
- Volume(大量):数据体量巨大,从 TB 级到 PB、EB 级。例如某电商每天产生数十亿条用户行为日志。
- Velocity(高速):数据产生和流动速度快,要求系统能够实时或近实时摄取与处理。例如传感器每秒上报数千次读数。
- Variety(多样):数据类型多样,包括结构化(关系表)、半结构化(JSON、XML、日志)、非结构化(文本、图像、音视频)。
- Veracity(真实性):数据质量参差不齐,存在噪声、缺失、不一致,需要数据清洗与质量治理。
- Value(价值):单条数据价值密度低,但整体蕴含巨大价值,需要通过分析挖掘才能释放。
· 体量大(Volume)、速度快(Velocity)、种类多(Variety)、真实性(Veracity)、价值高(Value)
· "体速种真价" —— 体量大、速度快、种类多、真实性、价值密度低但总量价值高
2.2 大数据处理流水线
无论采用何种具体技术,一个完整的大数据处理流水线通常包含五个核心环节:采集 → 存储 → 处理 → 分析 → 可视化。每一环节都有对应的代表性技术栈。
需要注意的是,这五个环节并非严格线性,实际系统中常常存在回环:例如分析结果可能触发新的采集任务,可视化发现的异常可能反过来驱动模型重训练。架构师的任务是把这些环节解耦但又串联起来,形成可扩展的数据中台。
三、大数据架构模式
大数据架构模式解决的核心问题是:如何同时满足海量数据的批处理(高吞吐、低成本)与实时处理(低延迟)需求。围绕这一问题,业界先后演化出 Lambda、Kappa 与湖仓一体三种主流架构。
3.1 Lambda 架构
Lambda 架构由 Nathan Marz 提出,其核心思想是"批处理 + 流处理"双链路并行,再通过服务层合并结果。它将系统分为三层:
- 批处理层(Batch Layer):用 Hadoop MapReduce / Spark 等对全量数据进行高吞吐、低成本的批处理,生成批处理视图(Batch View)。优点是结果精确,缺点是延迟高(小时~天级)。
- 速度层(Speed Layer):用 Storm / Flink / Spark Streaming 等对增量数据进行低延迟处理,生成实时视图(Realtime View)。优点是延迟低(秒~毫秒级),缺点是结果近似、资源消耗高。
- 服务层(Serving Layer):合并批处理视图与实时视图,对外提供低延迟的查询服务(如 HBase、Druid、Presto)。
✓ 优点
- 兼顾低延迟与结果精确,容错性强
- 批处理层可重算历史数据,结果可重现
- 架构成熟,工业界广泛落地
✗ 缺点
- 需维护两套代码(批+流),开发与运维成本高
- 两套逻辑结果可能不一致,调试困难
- 资源占用大,集群规模翻倍
3.2 Kappa 架构
Kappa 架构由 Jay Kreps(Kafka 作者)提出,主张用单一的流处理链路取代批+流双链路。其核心思想是:所有数据都视为流,批处理只是"有界流"的特例。当需要重算历史时,通过Kafka 消息回放(replay)重新消费历史数据即可,无需维护两套代码。
✓ 优点
- 只维护一套代码,开发运维成本低
- 结果一致性更好
- 架构简洁,易于演进
✗ 缺点
- 依赖消息队列的回放能力与长保留期
- 大范围历史重算开销大
- 对流处理引擎要求高(状态管理、Exactly-once)
3.3 湖仓一体(Lakehouse)
湖仓一体由 Databricks 提出,旨在融合数据湖的灵活性与数据仓库的管理能力。传统数据湖基于 HDFS/S3 存储原始数据,成本低、格式灵活,但缺乏事务、Schema 强约束、数据质量保障;数据仓库则相反,管理规范但成本高、扩展性差。湖仓一体通过在对象存储之上引入开放表格式(Open Table Format),实现 ACID 事务、Schema 演进、Time Travel 等能力。
三大主流开放表格式:
- Delta Lake(Databricks 主导):基于 Parquet,支持 ACID、Schema 演进、Time Travel、Z-Order 优化。
- Apache Iceberg(Netflix 主导):表格式规范,引擎解耦,支持隐藏分区、分区演进,与 Spark/Flink/Trino 等多引擎兼容。
- Apache Hudi(Uber 主导):强项在 Upsert/增量处理,适合 CDC 场景与近实时数仓。
3.4 三种架构对比
| 维度 | Lambda 架构 | Kappa 架构 | 湖仓一体(Lakehouse) |
|---|---|---|---|
| 处理链路 | 批 + 流双链路 | 单一流处理链路 | 存储计算分离,多引擎统一访问 |
| 代码维护 | 两套代码 | 一套代码 | 一套表格式 + 多引擎 |
| 历史重算 | 批处理层重算 | Kafka 回放 | Time Travel / 快照 |
| 事务支持 | 弱 | 依赖引擎 | ACID(开放表格式提供) |
| 实时性 | 高(速度层) | 很高 | 中高(支持流式写入) |
| 成本 | 高(双链路) | 中 | 低(对象存储) |
| 典型技术 | Hadoop + Storm + HBase | Kafka + Flink | S3 + Iceberg/Delta/Hudi + Spark |
| 适用场景 | 既有历史又有实时 | 以实时为主、回放可行 | 统一数据平台、ML/BI 一体 |
四、Hadoop 生态系统
Hadoop 是大数据时代的奠基性开源框架,由 Doug Cutting 以 Google 三大论文(GFS、MapReduce、Bigtable)为蓝本实现。它解决了海量数据的分布式存储与分布式计算问题,是软考大数据部分的高频考点。
4.1 HDFS(Hadoop 分布式文件系统)
HDFS 采用主从架构,将大文件切分为固定大小的块(Block,默认 128MB)分布式存储在多个节点上,并通过多副本(默认 3 副本)保障容错。
- NameNode:主节点,管理文件系统的元数据(命名空间、Block 到 DataNode 的映射、副本位置)。不存储实际数据,内存中维护。
- DataNode:从节点,存储实际的数据块,定期向 NameNode 发送心跳和块报告。
- Secondary NameNode:辅助节点,并非 NameNode 的热备,而是定期合并 fsImage 与 editLog,减轻 NameNode 重启压力。
Secondary NameNode ≠ NameNode 备份!它的作用是定期合并元数据镜像,降低 NameNode 启动时间。真正的 HA 备份通常通过HA NameNode(Active/Standby)+ JournalNode 实现。
· Block 大小:默认 128MB(Hadoop 2.x+),旧版 64MB
· 副本系数:默认 3,存放策略为"本地一份、同机架一份、跨机架一份"
· 写入流程:Client → NameNode(申请)→ DataNode(流水线复制)→ ACK
· 读取流程:Client → NameNode(获取块位置)→ 就近 DataNode 拉取
4.2 MapReduce
MapReduce 是一种分而治之的分布式计算模型,将任务分为 Map(映射)和 Reduce(归约)两个阶段,中间通过 Shuffle(洗牌)将相同 key 的数据分发到同一 Reduce 任务。
- Map 阶段:将输入分片(InputSplit)并行处理,输出 <key, value> 键值对。
- Shuffle 阶段:对 Map 输出按 key 分区、排序、分组,传输到 Reduce 端。Shuffle 是 MapReduce 性能瓶颈所在。
- Reduce 阶段:对相同 key 的所有 value 进行聚合计算,输出最终结果。
1. 输入分片 → 每个 Map 处理一行文本,按空格切分,每个单词输出 (word, 1)
2. Shuffle 按 word 分区排序,相同 word 的 1 聚合为 list:[1,1,...]
3. Reduce 对 list 求和,输出 (word, 总次数)
4. Shuffle 是性能瓶颈:大量磁盘 I/O 与网络传输,这是 Spark 取代 MapReduce 的关键原因
4.3 YARN(资源调度)
YARN(Yet Another Resource Negotiator)是 Hadoop 2.x 引入的资源管理与调度框架,将资源管理与作业调度解耦,使 Hadoop 不再只服务于 MapReduce,而是可以同时运行 Spark、Flink、Tez 等多种计算框架。
- ResourceManager:全局资源管理者,负责集群资源调度与作业调度。内含 Scheduler(调度器)和 ApplicationsManager(应用管理器)。
- NodeManager:每个节点上的代理,负责容器(Container)的生命周期与资源监控。
- ApplicationMaster:每个应用(如一个 Spark 作业)一个,负责申请资源、启动任务、容错重试。
- Container:资源抽象(CPU + 内存),任务运行的隔离环境。
1. FIFO Scheduler:先进先出,单队列,简单但公平性差
2. Capacity Scheduler(默认):多队列,按容量分配,保证小作业资源
3. Fair Scheduler:公平共享,所有作业平均分配资源,适合多用户
4.4 Hadoop 生态全景
Hadoop 已从最初的 HDFS + MapReduce 发展为庞大生态,覆盖数据接入、存储、计算、查询、协调、调度等各个环节。
| 组件 | 类别 | 核心作用 |
|---|---|---|
| HDFS | 存储 | 分布式文件系统,海量数据底座 |
| HBase | 存储 | 列式 NoSQL 数据库,基于 HDFS,支持随机读写 |
| MapReduce | 计算 | 分布式批处理框架 |
| YARN | 调度 | 集群资源管理与作业调度 |
| Hive | 查询 | SQL on Hadoop,将 SQL 翻译为 MapReduce/Tez 作业 |
| Pig | 查询 | 数据流脚本语言 Pig Latin,适合 ETL |
| Sqoop | 接入 | 关系数据库 ↔ HDFS/Hive/HBase 数据同步 |
| Flume | 接入 | 日志流式采集,写入 HDFS/Kafka |
| Kafka | 消息 | 分布式消息队列,削峰填谷、流处理入口 |
| ZooKeeper | 协调 | 分布式协调服务:配置/命名/锁/选举 |
| Oozie | 调度 | 工作流调度,编排 MapReduce/Hive 作业 |
| Impala/Presto | 查询 | MPP 引擎,交互式低延迟 SQL 查询 |
五、Spark 内存计算
Apache Spark 是继 MapReduce 之后的新一代分布式计算引擎,核心优势是基于内存的迭代计算,相比 MapReduce 的磁盘密集型 Shuffle,Spark 可将中间结果驻留内存,性能提升 10~100 倍,特别适合机器学习迭代算法与交互式查询。
5.1 Spark 架构
- Driver:运行用户 main 函数,创建 SparkContext,构建 DAG、调度 Task。
- Executor:工作节点上的 JVM 进程,负责运行 Task 并缓存数据(Cache)。
- Cluster Manager:集群资源管理,可选 Standalone、YARN、Mesos、K8s。
- Task:最小执行单元,一个 Stage 内被切分为多个 Task 并行执行。
5.2 Spark 核心概念
- RDD(Resilient Distributed Dataset):弹性分布式数据集,Spark 最底层抽象。只读、可分区、可容错(通过血缘 Lineage 重建)。具有 5 大特性:分区列表、计算函数、依赖列表、Partitioner、首选位置。
- DAG(有向无环图):RDD 之间的转换关系构成 DAG,Driver 将 DAG 划分为 Stage。
- Stage:以 Shuffle 为边界切分,Stage 内部全是窄依赖,可流水线并行。
- Task:Stage 内按分区切分的执行单元,一个分区一个 Task。
- 宽依赖 vs 窄依赖:窄依赖(map、filter)父子分区一对一;宽依赖(groupByKey、join)一对多,触发 Shuffle。
· RDD:底层 API,灵活但无优化
· DataFrame:带 Schema 的 RDD,类似关系表,有 Catalyst 优化器与 Tungsten 执行引擎
· Dataset:DataFrame + 强类型(Scala/Java),兼顾类型安全与性能
5.3 Spark 生态栈
- Spark Core:RDD、调度、内存管理、容错,所有模块的基础。
- Spark SQL:SQL 与 DataFrame/Dataset 接口,可访问 Hive、Parquet、JSON 等。
- Spark Streaming:微批流处理(Micro-Batch),将流切分为小批次处理。
- MLlib:机器学习库,含分类、回归、聚类、协同过滤、推荐等算法。
- GraphX:图计算库,支持 Pregel API 与图算法(PageRank、连通分量等)。
5.4 Spark vs MapReduce
| 维度 | MapReduce | Spark |
|---|---|---|
| 计算模式 | 批处理 | 批 + 微批流 + 机器学习 + 图 |
| 中间结果 | 落磁盘 | 驻留内存(可落盘) |
| 性能 | 慢(磁盘 I/O) | 快 10~100 倍 |
| 编程模型 | Map/Reduce 两阶段 | RDD/DAG 多阶段 |
| 迭代计算 | 不友好(每次重读磁盘) | 友好(内存缓存) |
| 容错 | Task 重试 | 血缘 Lineage 重建 |
| 延迟 | 分钟~小时 | 秒~分钟 |
| 资源 | 低 | 高(内存大) |
六、Flink 流处理
Apache Flink 是真正的流处理引擎(Native Streaming),与 Spark Streaming 的微批(Micro-Batch)模式形成鲜明对比。Flink 把每条事件当作流的一个元素逐条处理,延迟可低至毫秒级,是当前低延迟流处理的事实标准。
6.1 真流 vs 微批
- Spark Streaming(微批):将流按时间窗口(如 1 秒)切分为小批次 RDD,每个批次内部仍是批处理。延迟取决于批次间隔(秒级)。
- Flink(真流):每条事件到来即处理,无需攒批,延迟毫秒级。批处理被视为"有界流"的特例,统一在流处理模型下。
6.2 Flink 架构
- JobManager:主节点,负责作业调度、检查点协调、故障恢复。内含 ResourceManager、Dispatcher、JobMaster。
- TaskManager:工作节点,执行 Task,管理 Slot(资源槽)、状态(State)与网络缓冲。
- Slot:TaskManager 内的资源单位,一个 Slot 可运行一个 Task 子任务。
6.3 关键特性
- 有状态流处理(Stateful Streaming):Flink 内置状态管理(ValueState、ListState、MapState),支持算子状态与键控状态,状态可持久化到 RocksDB。
- 事件时间处理(Event Time):使用事件本身的时间戳而非系统处理时间,保证结果确定性。
- 水位线(Watermark):用于处理乱序事件,表示"小于该时间的事件已全部到达",触发窗口计算。
- Exactly-Once 语义:通过 Checkpoint(基于 Chandy-Lamport 算法)+ 两阶段提交(2PC)Sink 实现"精确一次"语义。
- Checkpoint vs Savepoint:Checkpoint 自动周期性触发用于故障恢复;Savepoint 手动触发用于版本升级与作业迁移。
Flink 端到端 Exactly-Once 需要三件套配合:① Source 支持回放(如 Kafka offset)② Flink Checkpoint ③ Sink 支持两阶段提交(如 Kafka 事务写、JDBC 事务)。缺一则只能 At-Least-Once。
6.4 Spark Streaming vs Flink
| 维度 | Spark Streaming | Flink |
|---|---|---|
| 处理模型 | 微批(Micro-Batch) | 真流(Native Streaming) |
| 延迟 | 秒级 | 毫秒级 |
| 时间语义 | 处理时间为主 | 事件时间 + Watermark |
| 状态管理 | 较弱 | 强大(内置 State) |
| 语义保障 | At-Least-Once | Exactly-Once(端到端) |
| 批处理 | 原生支持 | 有界流(DataSet 已弃用,统一 DataStream) |
| 生态 | SQL/MLlib/GraphX 完善 | SQL/CEP/Table 完善 |
| 适用场景 | 批为主、流为辅 | 低延迟流处理为主 |
七、数据仓库与数据湖
7.1 数据仓库分层架构
数据仓库通常采用分层建模,将原始数据逐步加工为可分析指标,典型四层架构:
- ODS(Operational Data Store,操作数据层):原始数据贴源层,结构基本与源系统一致,保留历史快照。
- DWD(Data Warehouse Detail,明细数据层):清洗、标准化、维度退化后的明细事实表。
- DWS(Data Warehouse Summary,汇总数据层):按主题、维度轻度汇总,如用户日汇总、商品日汇总。
- ADS(Application Data Service,应用数据层):面向应用的指标结果,直接对接 BI 报表、API。
· 复用:DWD/DWS 一次加工多处使用,避免烟囱式开发
· 解耦:源系统变更不直接影响应用层
· 血缘清晰:便于数据治理与影响分析
7.2 星型模型 vs 雪花模型
数仓维度建模有两种典型模式:
- 星型模型(Star Schema):一张事实表居中,多张维度表直接围绕,维度表不再规范化。结构简单、查询性能好,是数仓首选。
- 雪花模型(Snowflake Schema):维度表进一步规范化拆分为多张子表,呈雪花状。节省存储但 JOIN 多、性能差。
| 维度 | 星型模型 | 雪花模型 |
|---|---|---|
| 维度表规范化 | 否(反范式) | 是(范式拆分) |
| JOIN 数量 | 少 | 多 |
| 查询性能 | 好 | 较差 |
| 存储 | 冗余多 | 节省 |
| 建模复杂度 | 简单 | 复杂 |
| 适用场景 | OLAP 数仓首选 | 关系强、需精细管理 |
7.3 缓慢变化维(SCD)
维度数据会随时间变化(如用户迁居、商品改名),如何记录历史变化?这就是缓慢变化维(Slowly Changing Dimensions)问题,常见三类处理方式:
- SCD Type 1(覆盖):直接用新值覆盖旧值,不保留历史。适用于不关心历史的字段(如修正错误)。
- SCD Type 2(拉链表):新增一行记录新值,并标记生效时间与失效时间,保留完整历史。是数仓最常用方式。
- SCD Type 3(新增列):在原行新增"旧值"列,仅保留上一版本。适用于只需对比前后两次的场景。
· Type 1 = 覆盖(无历史)
· Type 2 = 拉链表(全历史,最常用)
· Type 3 = 加列(仅上一版)
拉链表关键字段:start_date、end_date、is_current
7.4 数据湖
数据湖(Data Lake)是一个集中式存储库,可以原生格式存储任意类型、任意规模的数据。与数据仓库的"先建模后入库"不同,数据湖主张"先入库后建模"(Schema-on-Read),灵活但易沦为"数据沼泽"。前述湖仓一体正是为解决数据湖管理薄弱而生的演进。
| 维度 | 数据仓库 | 数据湖 |
|---|---|---|
| 数据类型 | 结构化为主 | 结构化/半/非结构化 |
| Schema | Schema-on-Write(先建表) | Schema-on-Read(读取时解析) |
| 处理者 | 业务分析师 | 数据科学家/工程师 |
| 成本 | 高 | 低 |
| 成熟度 | 成熟 | 易成数据沼泽 |
八、NoSQL 数据库在大数据中的应用
8.1 NoSQL 四大类型回顾
| 类型 | 代表 | 数据模型 | 典型场景 |
|---|---|---|---|
| 键值(KV) | Redis、Memcached | Key-Value | 缓存、会话、计数器 |
| 文档(Document) | MongoDB、CouchDB | JSON 文档 | 内容管理、用户画像 |
| 列族(Column-Family) | HBase、Cassandra | 稀疏列族 | 大数据宽表、时序 |
| 图(Graph) | Neo4j、JanusGraph | 节点+边 | 社交、推荐、风控 |
8.2 列族数据库与 HBase
列族数据库是大数据场景的主力,能够支撑数十亿行 × 百万列的稀疏宽表。HBase 基于 HDFS,提供强一致性随机读写,是 Hadoop 生态的核心 NoSQL。
HBase 核心组件:
- HMaster:主节点,负责 Region 分配、负载均衡、Schema 变更。
- RegionServer:数据服务节点,管理多个 Region,处理读写请求。
- Region:表按 RowKey 范围切分的分片,是 HBase 水平扩展的最小单位,可达自动分裂。
- HLog:预写日志(WAL),写入前先落盘,保障故障恢复。
- MemStore:内存写缓冲,写满后 flush 为 HFile。
- HFile:HDFS 上的列式存储文件,LSM-Tree 结构。
- Store:一个 Region 内每个列族对应一个 Store。
1. Client 通过 ZooKeeper 定位 RegionServer
2. 写 HLog(WAL)落盘保障持久性
3. 写 MemStore(内存)
4. 返回成功
5. MemStore 写满后 flush 为 HFile
6. 多个 HFile 定期 Compaction 合并
读取流程:MemStore + BlockCache + HFile 三处合并扫描
8.3 各类 NoSQL 在大数据中的适用场景
- 键值(Redis):实时推荐的特征缓存、热门商品榜单、限流计数。
- 文档(MongoDB):用户画像标签、商品元数据、内容管理(灵活 Schema 适合频繁变字段)。
- 列族(HBase/Cassandra):海量订单流水、监控时序数据、风控行为日志。
- 图(Neo4j):社交关系、欺诈团伙挖掘、知识图谱、推荐二跳关系。
九、消息队列与流处理
9.1 Kafka 架构
Apache Kafka 是分布式发布-订阅消息系统,以高吞吐、可持久化、可水平扩展著称,是大数据管道的"主动脉"。
- Producer:消息生产者,按 key 将消息写入对应 Partition。
- Consumer:消息消费者,从 Partition 拉取消息。
- Broker:Kafka 服务节点,集群由多个 Broker 组成。
- Topic:消息主题,逻辑分类。
- Partition:Topic 的分片,是并行度单位,每个 Partition 内消息有序。
- Consumer Group:消费者组,组内消费者分担 Partition,一条消息只被组内一个消费者消费。
- Offset:消费者在 Partition 中的消费位移。
- Replica:副本,分为 Leader(读写)与 Follower(同步)。
9.2 Kafka vs RabbitMQ
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 定位 | 分布式日志/流平台 | 传统消息中间件 |
| 吞吐 | 极高(百万/秒) | 中(万级/秒) |
| 延迟 | 毫秒~秒 | 微秒~毫秒 |
| 持久化 | 磁盘顺序写,可长期保留 | 可持久化,但消费后删除 |
| 消费模型 | 拉取(Pull),按 Offset 回放 | 推送(Push),消费即删除 |
| 顺序性 | Partition 内有序 | 队列内有序 |
| 典型场景 | 日志、流处理、事件溯源 | 业务解耦、异步任务、RPC |
9.3 流处理模式
- 无状态转换:map、filter,每条独立处理,无需保存历史。
- 有状态聚合:窗口聚合(求和、计数)、Keyed State,需要保存中间状态。
- 窗口(Window):滚动窗口(不重叠)、滑动窗口(可重叠)、会话窗口(按活动间隙)。流处理核心概念。
- 双流 Join:Interval Join、Window Join,将两条流按时间窗口关联。
- 复杂事件处理(CEP):检测事件序列模式,如"3 次失败登录后 1 次成功"。
· 顺序写磁盘:追加写入,磁盘顺序写接近内存速度
· 零拷贝(Zero Copy):sendfile 系统调用,数据不经用户态
· 分区并行:Topic 多 Partition,水平扩展
· 批量压缩:Producer 端批量 + 压缩(snappy/lz4)
十、机器学习架构
10.1 机器学习流程
一个完整的机器学习项目并非只有"训练模型"一步,而是覆盖从数据到部署到监控的完整闭环,通常包含六个阶段:
- 数据采集:从业务系统、日志、第三方 API 收集原始数据。
- 特征工程:清洗、缺失值处理、特征提取、特征变换、特征选择。
- 模型训练:选择算法、划分训练/验证集、超参数调优、交叉验证。
- 模型评估:在测试集上评估准确率、精确率、召回率、AUC、F1 等指标。
- 模型部署:将模型打包为服务(REST/gRPC),上线提供推理。
- 模型监控:监控预测分布、数据漂移、业务指标,触发再训练。
10.2 特征工程与特征存储
特征工程是机器学习中最耗工时的环节,业界常说"数据和特征决定了机器学习的上限,而算法只是逼近这个上限"。
- 特征存储(Feature Store):集中管理特征定义、计算逻辑与特征值,统一训练与推理的特征口径,避免"训练-在线特征不一致"问题。代表实现:Feast、Tecton、Hopsworks。
- 离线特征:面向训练,从数据仓库批量计算,规模大、延迟容忍。
- 在线特征:面向推理,存于 Redis/Feature Store,毫秒级读取。
若训练用离线特征计算逻辑、在线用另一套代码生成特征,会导致模型上线效果骤降。Feature Store 通过单一特征定义 + 双写统一两套口径,是解决此问题的关键架构。
10.3 模型训练与分布式训练
当模型规模(如大模型 LLM)或数据量超出单机内存/算力时,需要分布式训练。两种基本并行策略:
- 数据并行(Data Parallel):每个 worker 持有完整模型副本,各处理不同数据分片,反向传播后通过 AllReduce 同步梯度。最常用。
- 模型并行(Model Parallel):将模型本身切分到多个 worker,每个 worker 只持有部分参数。适用于单机装不下的超大模型。
- 流水线并行(Pipeline Parallel):将模型按层切分到不同设备,形成流水线。常与数据并行组合用于大模型训练。
参数服务器(Parameter Server)架构:经典的分布式训练架构,由 Server 节点(存参数)与 Worker 节点(算梯度)组成,Worker pull 参数、push 梯度,Server 异步或同步更新。
10.4 模型推理服务
模型训练完成后,需要部署为推理服务(Model Serving)对外提供预测能力。两种推理模式:
- 批量推理(Batch Inference):离线对一批样本批量预测,适合非实时场景(如日级推荐候选生成)。
- 实时推理(Real-time Inference):在线服务,单条请求毫秒级响应,适合广告竞价、风控实时拦截。
主流模型服务框架:
- TensorFlow Serving:Google 开源,支持 TF 模型,gRPC/REST 接口,模型版本管理。
- TorchServe:PyTorch 官方服务框架,支持多模型、批处理。
- NVIDIA Triton:支持多框架(TF/PyTorch/ONNX/TensorRT),GPU 加速,动态批处理。
- ONNX Runtime:跨框架统一推理,支持 CPU/GPU 优化。
十一、MLOps 与模型生命周期管理
MLOps(Machine Learning Operations)是将 DevOps 思想应用于机器学习的工程实践,目标是让模型从实验到生产的全流程自动化、可监控、可复现。MLOps 解决的核心痛点是:模型实验环境与生产环境脱节、模型上线慢、上线后效果衰退无人感知。
11.1 MLOps 核心能力
- ML 流水线自动化(Pipeline Automation):将数据准备、训练、评估、部署编排为可重复执行的流水线,支持定时触发与事件触发。
- 模型版本管理(Model Versioning):像代码版本控制一样管理模型,记录每个版本的数据、代码、超参、指标。代表:MLflow、DVC。
- 模型注册中心(Model Registry):集中存储模型及其元数据,管理模型状态(Staging/Production/Archived)。
- 模型监控(Model Monitoring):上线后持续监控预测分布、特征分布、业务指标。
- 持续训练(CT, Continuous Training):当数据漂移或定时触发时自动重训练,是 MLOps 区别于 DevOps 的关键。
11.2 数据漂移与概念漂移
模型上线后效果会随时间衰退,根因主要有两类"漂移":
- 数据漂移(Data Drift):输入特征分布变化,但 X→Y 关系不变。例如用户年龄结构变化。可通过 PSI、KS 检验、KL 散度检测。
- 概念漂移(Concept Drift):X→Y 关系本身变化,即"规则变了"。例如疫情导致消费偏好剧变。检测更难,需结合业务指标。
· 轻微漂移 → 监控告警,暂不重训
· 显著漂移 → 触发自动重训练(CT)
· 概念漂移 → 需引入新特征或更换模型,单纯重训无效
十二、AI 推理优化
模型在训练阶段追求精度,在推理阶段则追求低延迟、高吞吐、低资源消耗。AI 推理优化是让模型在生产环境中"跑得快又省钱"的关键技术,主要手段有三类:量化、剪枝、知识蒸馏。
12.1 模型量化(Quantization)
量化是将模型权重与激活值从高精度浮点(FP32)转换为低精度整数(INT8/INT4)的过程,可大幅降低内存占用、提升推理速度,是工业部署最常用的优化手段。
- 训练后量化(Post-Training Quantization, PTQ):训练完成后用少量校准数据统计激活范围,直接转换为 INT8。简单快速,但可能掉点。
- 量化感知训练(Quantization-Aware Training, QAT):训练时模拟量化误差,让模型适应低精度。精度更高但训练成本大。
12.2 模型剪枝(Pruning)
剪枝是删除模型中冗余、不重要的权重或神经元,减小模型体积。
- 非结构化剪枝(Unstructured Pruning):逐个权重置零,产生稀疏矩阵。压缩率高但通用硬件难以加速,需专用稀疏库。
- 结构化剪枝(Structured Pruning):按通道/层/块整体删除,硬件友好,可直接加速。工业更常用。
12.3 知识蒸馏(Knowledge Distillation)
知识蒸馏用一个大而强的教师模型(Teacher)指导训练一个小而快的学生模型(Student),让学生学到教师的"暗知识"(soft label 的概率分布),在保持接近教师精度的同时大幅降低推理成本。
| 技术 | 原理 | 加速效果 | 精度损失 | 典型工具 |
|---|---|---|---|---|
| 量化 | FP32→INT8/INT4 | 2-4 倍 | 轻微 | TensorRT/TFLite |
| 剪枝 | 删除冗余权重 | 1.5-3 倍 | 可控 | TorchPruner |
| 蒸馏 | 大模型教小模型 | 取决于学生 | 较小 | HuggingFace/自研 |
| 算子融合 | 合并相邻算子 | 1.2-2 倍 | 无损 | TensorRT/XLA |
工业界通常组合使用:先蒸馏得到小模型 → 再量化为 INT8 → 再结构化剪枝 → 最后算子融合。最终在 GPU/CPU 上可获得数倍加速且精度损失可控。
十三、大数据与 AI 实战案例:推荐系统架构
推荐系统是大数据与 AI 融合的典型场景,几乎串联了本文所述的所有技术:数据采集、消息队列、流处理、特征工程、模型训练、模型推理、A/B 测试、监控反馈。下面以电商推荐为例,剖析完整架构。
13.1 整体架构
现代推荐系统普遍采用"召回 → 排序 → 重排"三段式架构,兼顾候选覆盖度与排序精度,同时满足实时性要求。
13.2 关键环节详解
- 数据采集:通过埋点 SDK 采集用户点击、加购、停留、购买行为,经 Kafka 进入实时与离线两条链路。商品元数据来自业务库(Binlog CDC 同步)。
- 特征工程:用户特征(性别、年龄段、历史偏好)、商品特征(类目、价格、热度)、上下文特征(时间、地点、设备)。特征统一存入 Feature Store,离线训练与在线推理共用。
- 召回(Recall):从百万级商品中初筛千级候选。多路并行:协同过滤(User-CF/Item-CF)、向量化召回(双塔模型 + Faiss 近邻检索)、图召回(基于知识图谱的 GNN)、热门召回、标签召回。
- 排序(Ranking):对千级候选精排打分,预测 CTR/CVR。常用模型:LR、GBDT、DeepFM、Wide&Deep、DIN(深度兴趣网络)。融合稠密特征与 embedding。
- 重排(Re-ranking):考虑多样性、新颖性、业务规则(去重、打散、强插广告),避免推荐同质化。
- 实时服务:API 网关接收请求,并行拉取在线特征(Redis)+ 调用召回/排序模型(Triton),毫秒级返回推荐列表。
- A/B 测试:按用户 ID 哈希分桶,对照组用旧模型、实验组用新模型,对比 CTR、GMV、留存等指标,统计显著性通过后全量上线。
- 监控反馈:实时监控推荐 CTR、曝光分布、模型预测分布,发现数据漂移即触发离线重训练,形成闭环。
1. 离线-在线一致:特征计算逻辑必须训练与推理同源,否则效果骤降
2. 召回求宽、排序求准:召回多路并行保证覆盖,排序精打细算保证质量
3. 实时性分级:用户实时行为(最近点击)秒级更新,长期偏好天级更新
4. 灰度发布:新模型先小流量灰度,验证无损再放量
十四、软考考点总结与真题
14.1 核心考点速查表
| 知识域 | 核心考点 | 考频 |
|---|---|---|
| 大数据基础 | 5V 特征、处理流水线 | ★★★★★ |
| 架构模式 | Lambda/Kappa/湖仓一体对比 | ★★★★ |
| Hadoop | HDFS 三角色、MapReduce 流程、YARN 调度 | ★★★★★ |
| Spark | RDD/DAG/Stage、与 MR 对比 | ★★★★ |
| Flink | 真流 vs 微批、Exactly-Once、Watermark | ★★★ |
| 数仓 | ODS/DWD/DWS/ADS、星型/雪花、SCD | ★★★★ |
| NoSQL | 四大类型、HBase 架构 | ★★★ |
| Kafka | 架构组件、与 RabbitMQ 对比 | ★★★ |
| ML/MLOps | ML 流程、特征存储、模型监控 | ★★★ |
| 推理优化 | 量化/剪枝/蒸馏区分 | ★★ |
14.2 软考真题精选
真题一(大数据特征)
大数据的基本特征通常用 5V 来描述,以下不属于 5V 的是( )。
A. Volume B. Velocity C. Visualization D. Value
解析:大数据 5V 为 Volume(大量)、Velocity(高速)、Variety(多样)、Veracity(真实性)、Value(价值)。Visualization(可视化)是大数据的应用环节,不属于 5V 特征。注意 Veracity 与 Visualization 的区分。
真题二(Hadoop 组件职责)
在 HDFS 中,关于 Secondary NameNode 的描述,正确的是( )。
A. 它是 NameNode 的热备份,故障时自动接管
B. 它负责合并 fsImage 与 editLog,减轻 NameNode 重启压力
C. 它存储实际数据块,处理客户端读写请求
D. 它负责集群资源调度与作业分配
解析:Secondary NameNode 的核心职责是定期合并元数据镜像 fsImage 与编辑日志 editLog,降低 NameNode 重启时间。它并非 NameNode 的热备(A 错),不存储数据块(C 错,那是 DataNode),不负责资源调度(D 错,那是 YARN 的 ResourceManager)。真正的 HA 备份通过 Active/Standby NameNode + JournalNode 实现。
真题三(架构模式辨析)
关于 Lambda 架构与 Kappa 架构,下列说法错误的是( )。
A. Lambda 架构同时维护批处理层与速度层,需两套代码
B. Kappa 架构通过消息队列回放历史数据实现重算,仅维护一套代码
C. Lambda 架构的批处理层结果精确,速度层结果近似
D. Kappa 架构完全不需要任何批处理,所有计算都是事件级实时处理
解析:Kappa 架构主张用统一的流处理取代"批+流双链路",但并非"完全不需要批处理"。在 Kappa 中,批处理被视为"有界流"的特例,历史重算通过 Kafka 回放批量消费历史数据完成,本质仍是流处理模型对批处理的统一。A、B、C 描述均正确。Kappa 的核心是"一套代码 + 消息回放",而非"没有批处理"。
14.3 记忆口诀
· 5V:"体速种真价"——大量、高速、多样、真实、价值
· 流水线:"采存处分可"——采集、存储、处理、分析、可视化
· Lambda 三层:"批速服"——Batch/Speed/Serving
· HDFS 三角:"名数副"——NameNode/DataNode/SecondaryNameNode(副非备份!)
· MapReduce:"Map 洗 Reduce"——映射、Shuffle、归约
· YARN:"RM+NM+AM"——ResourceManager/NodeManager/ApplicationMaster
· Spark 核心:"RDD 造 DAG,宽依赖切 Stage,分区定 Task"
· Flink 特性:"真流、事件时、水位线、精确一次"
· 数仓分层:"ODW A"——ODS/DWD/DWS/ADS
· SCD:"一覆盖、二拉链、三加列"
· HBase 写入:"WAL→MemStore→Flush→HFile"
· 推理优化:"量剪蒸"——量化、剪枝、蒸馏
· 推荐三段:"召排重"——召回、排序、重排
14.4 总结
大数据与 AI 架构是系统架构设计师考试中覆盖面最广、技术更新最快的领域之一。从底层的 Hadoop 生态,到内存计算的 Spark、流处理的 Flink,再到数据仓库与数据湖的融合,直至机器学习与 MLOps 的工程化落地,构成了一条完整的"数据→智能"价值链。
1. 建立体系:不要孤立记忆每个组件,要理解其在数据流水线中的位置
2. 抓对比:Lambda vs Kappa、Spark vs Flink、星型 vs 雪花、数据仓库 vs 数据湖、量化 vs 剪枝 vs 蒸馏——对比是软考命题最爱
3. 记组件职责:HDFS/YARN/HBase/Kafka 的核心组件角色是高频选择/判断题
4. 理解架构演进动机:每项技术为何出现、解决了什么问题、代价是什么,比死记更重要
5. 结合案例:推荐系统案例串联了全部知识点,是论述题的好素材
架构的本质是权衡。没有最好的架构,只有最适合业务场景、团队规模、成本约束的架构。大数据与 AI 架构的演进史,正是工程师们在"延迟 vs 吞吐、成本 vs 性能、灵活 vs 规范、实时 vs 精确"之间不断权衡的历史。掌握这些权衡的思想,远比记忆具体组件的配置参数更有价值——这也是系统架构设计师考试的真正旨归。