Azuxa's Blog

大数据技术 系统性复习文档

系统性复习文档:涵盖大数据概述、数据存储与管理、分布式文件系统 HDFS、MapReduce 等核心章节。

技术 | 发布日期 2026/7/22

大数据技术 系统性复习文档


目录


第1章 大数据概述

1.1 大数据时代

1.1.1 第三次信息化浪潮

  • 根据IBM前CEO郭士纳的观点,IT领域每隔15年迎来一次重大变革
  • 三次信息化浪潮:
    • 第一次:以计算机为代表
    • 第二次:以互联网为代表
    • 第三次:以物联网为代表(大数据时代)

1.1.2 信息科技为大数据时代提供技术支撑

  • 存储设备容量不断增加,存储价格随时间下降
  • CPU处理能力大幅提升
    • 过去40多年:CPU从10MHz提高到4.6GHz
    • 遵循**“摩尔定律”**:芯片集成元件数量约每18个月翻一番,性能翻倍,价格下降一半
    • CPU制造工艺、晶体管数目、运行频率、核心数都在提升
  • 网络带宽不断增加
    • 1977年第一个光纤通信系统投入商用,速率45Mbit/s
    • 截至2022年底:我国互联网宽带接入端口10.65亿个,光纤占比95.7%,光缆线路5791万千米
    • 4G基站590万个(全球第一)
    • 截至2023年9月:5G基站318.9万个,5G用户7.37亿户

1.1.3 数据产生方式的变革

三个阶段:

  1. 运营式系统阶段:数据被动产生,伴随运营活动存入数据库
  2. 用户原创内容阶段:数据主动产生(Web 2.0时代),智能手机加速内容产生
  3. 感知式系统阶段:传感器广泛使用,数据自动采集,促成大数据产生

1.1.4 大数据的发展历程

  • 大数据发展经历三个阶段

1.2 大数据概念

1.2.1 数据量大(Volume)

  • 人类数据每年以**50%**速度增长(大数据摩尔定律)
  • 每两年增加一倍,最近两年产生数据 = 之前全部数据之和
  • 预测:2025年全球数据量175ZB,2030年达2500ZB
  • 中国:2025年预计48.6ZB,占全球27.8%

1.2.2 数据类型繁多(Variety)

  • 10%结构化数据:存于数据库中
  • 90%非结构化数据:与人类信息密切相关

1.2.3 处理速度快(Velocity)

  • 数据生成到消耗的时间窗口非常小
  • 1秒定律:与传统数据挖掘技术本质不同

1.2.4 价值密度低(Value)

  • 以视频为例:连续监控中可能仅一两秒有用数据,但商业价值很高

4V特征:Volume(大量)、Variety(多样)、Velocity(快速)、Value(价值)

1.3 大数据的影响

科学研究范式的演变(第四范式)

图灵奖得主Jim Gray总结:实验→理论→计算→数据(第四范式)

思维方式变革

  • 全样而非抽样
  • 效率而非精确
  • 相关而非因果

社会发展影响

  • 大数据决策成为新的决策方式
  • 促进信息技术与各行业深度融合
  • 推动新技术和新应用不断涌现
  • 就业:数据科学家成为热门职业
  • 教育:改变高校信息技术专业教学科研体制

1.4 大数据的应用

应用三个层次

  1. 描述性分析:总结抽取信息,分析”发生了什么”
    • 例:DOMO公司从企业数据中提取信息推送
  2. 预测性分析:分析关联关系和发展模式,预测趋势
    • 例:David Rothschild用大数据预测奥斯卡(21/24准确)
  3. 指导性分析:分析决策后果,指导优化决策
    • 例:无人驾驶汽车实时决策

典型应用实例

  • 《纸牌屋》:大数据分析决定选角和内容
  • 谷歌流感趋势:跟踪搜索词判断流感情况

1.5 大数据关键技术

  • 大数据技术不同层面及功能
  • 两大核心技术

1.6 大数据计算模式

计算模式代表产品
批处理MapReduce、Spark
流计算Storm、Flink、Spark Streaming
交互式查询Impala
图计算Pregel、GraphX

1.7 大数据产业

  • 大数据产业:一切与支撑大数据组织管理和价值发现相关的企业经济活动的集合

1.8 大数据与云计算、物联网的关系

1.8.1 云计算

  • 定义:通过网络提供可伸缩的、廉价的分布式计算能力
  • 三种服务模式
    • IaaS(Infrastructure as a Service):将基础设施(计算+存储)作为服务出租
    • PaaS(Platform as a Service):类似IaaS但包含OS和特定应用必需服务
    • SaaS(Software as a Service):集中部署软件,按使用计量收费
  • 三种部署模式:公有云、私有云、混合云
  • 关键技术:虚拟化、分布式存储、分布式计算、多租户
  • 数据中心:云计算的重要载体,提供计算、存储、带宽等资源

1.8.2 物联网

  • 定义:物物相连的互联网,互联网的延伸
  • 技术架构:感知层→网络层→应用层
  • 关键技术:识别和感知技术(二维码、RFID、传感器)、网络通信技术、数据挖掘与融合技术
  • 应用领域:智能交通、智慧医疗、智能家居、环保监测、智能安防、智能物流、智能电网、智慧农业、智能工业
  • 产业链:核心感应器件→感知层末端设备→网络→软件与方案→系统集成→运营服务

1.8.3 三者关系

  • 云计算、大数据、物联网相辅相成,既有联系又有区别
  • 大数据离不开云计算的支撑
  • 物联网是大数据的重要数据来源
  • 云计算为物联网和大数据提供计算能力

第2章 数据存储与管理

2.1 数据的概念

文件系统

  • 操作系统中负责管理和存储文件信息的软件机构
  • 由三部分组成:文件系统接口、软件集合、对象及属性
  • 负责为用户建立、存入、读出、修改、转储、撤销文件

数据库

  • 按一定规则组织、可共享、冗余少、与程序分开的数据集合
  • 特点:按一定方式存储、多用户共享、冗余度小、与应用程序独立

数据库发展历史

  1. 层次数据库(树形结构):简单但关系不够灵活
  2. 网状数据库:多对多联系,复杂难用
  3. 关系数据库:二维表存储,简单清晰易理解

关系数据库

  • 用二维表(行和列)组织数据
  • 例:学生信息表

数据仓库

  • 定义:面向主题的、集成的、相对稳定的、反映历史变化的数据集合,用于支持管理决策
  • 数据仓库体系架构

并行数据库

  • 在无共享体系结构中进行数据操作的数据库系统
  • 采用关系数据模型,支持SQL
  • 两个关键技术:关系表的水平划分、SQL查询的分区执行
  • 目标:高性能和高可用性
  • 缺点:数据转移代价昂贵

2.2 大数据时代数据存储与管理技术

分布式文件系统

  • 通过网络实现文件在多台主机上分布式存储
  • GFS:谷歌开发的分布式文件系统
  • HDFS:GFS的开源实现,Hadoop两大核心之一

NewSQL和NoSQL数据库

NewSQL数据库

  • 新的可扩展、高性能数据库
  • 兼具海量数据存储能力和传统数据库的ACID、SQL支持
  • 代表产品:Spanner、VoltDB、Clustrix等

NoSQL数据库

  • 非关系型数据库的统称
  • 没有固定表结构,无连接操作,不严格遵守ACID
  • 灵活的水平可扩展性,支持海量数据
  • 数据模型:键/值、列族、文档等非关系模型

大数据引发的数据库架构变革

  • 从单机关系数据库 → 并行数据库 → 分布式数据库 → NoSQL数据库

云数据库

  • 云服务提供商在公有云中托管数据库
  • 满足企业动态数据存储需求和中小企业低成本存储需求
  • 数据量每年60%速度增加

数据湖

  • 定义:企业中全量数据的单一存储
  • 包含:结构化、半结构化、非结构化、二进制数据
  • 可构建在本地数据中心或云上
  • 本质:数据存储架构 + 数据处理工具
  • 存储底座:通常用对象存储(如Amazon S3)
  • 两类工具
    1. 数据移动和管理工具(搬到湖里)
    2. 数据分析工具(从湖里淘金)

数据湖与数据仓库的区别

  • 数据仓库:结构化数据,模式固定,适合报表和分析
  • 数据湖:全量数据,模式灵活,适合探索和机器学习

湖仓一体

  • 打通数据仓库和数据湖的新型开放式架构
  • 融合数据仓库的高性能管理能力与数据湖的灵活性
  • 特性:存算分离、事务支持、开放性、数据治理、支持多种数据类型、BI支持
  • 核心能力:湖和仓的数据/元数据无缝打通、自由流动

2.3 Hadoop

Hadoop特性

  • 能够对大量数据进行分布式处理的软件框架
  • 可靠、高效、可伸缩
  • 特点:高可靠性、高效性、高可扩展性、高容错性、成本低、运行Linux、支持多种语言

Hadoop生态系统

  • 核心组件:HDFS(存储)+ MapReduce(计算)
  • 扩展组件:HBase、Hive、Pig、ZooKeeper、Sqoop、Flume等

2.4 HDFS

  • 设计目标:兼容廉价硬件、流数据读写、大数据集、简单文件模型、跨平台兼容
  • 局限性:不适合低延迟访问、无法高效存储小文件、不支持多用户写入和任意修改
  • 体系结构:一个NameNode + 若干DataNode

2.5 HBase

  • 开源的分布式存储系统,BigTable的开源实现
  • 特点:高可靠、高性能、面向列、可伸缩
  • BigTable可以扩展到PB级别数据、上千台机器

第3章 分布式文件系统HDFS(上)

3.1 分布式文件系统

3.1.1 计算机集群结构

  • 分布式文件系统把文件分布存储到多个计算机节点上
  • 节点由普通硬件构成,降低硬件开销
  • 两类节点:
    • 主节点(Master/NameNode)
    • 从节点(Slave/DataNode)

3.2 HDFS简介

  • 设计目标
    • 兼容廉价硬件设备
    • 流数据读写(一次写入、多次读取)
    • 大数据集(GB、TB、PB级)
    • 简单的文件模型
    • 强大的跨平台兼容性
  • 局限性
    • 不适合低延迟数据访问
    • 无法高效存储大量小文件(NameNode内存受限)
    • 不支持多用户写入及任意修改文件

3.3 HDFS相关概念

3.3.1 块(Block)

  • 默认块大小128MB,远大于普通文件系统(通常4KB)
  • 文件被分成多个块,以块为存储单位
  • 块的优点:
    1. 支持大规模文件存储
    2. 简化系统设计(固定大小简化元数据管理)
    3. 适合数据备份(利于容错和数据可用性)

3.3.2 名称节点和数据节点

NameNode(名称节点)

  • 管理分布式文件系统的命名空间(Namespace)
  • 保存两个核心数据结构:
    • FsImage:维护文件系统树及所有文件/文件夹的元数据
    • EditLog:记录所有文件操作(创建、删除、重命名等)
  • 记录每个文件中各块所在DataNode的位置信息

SecondaryNameNode(第二名称节点)

  • 解决EditLog不断变大的问题
  • 保存NameNode的元数据备份
  • 减少NameNode重启时间
  • 通常单独运行在一台机器上
  • 工作流程:定期下载FsImage和EditLog→合并→上传新的FsImage

DataNode(数据节点)

  • 工作节点,负责数据的存储和读取
  • 根据客户端或NameNode调度进行数据存储和检索
  • 定期向NameNode发送所存储块的列表
  • 数据保存在各节点本地Linux文件系统中

3.4 HDFS体系结构

3.4.1 体系结构概述

  • 主从(Master/Slave)结构
  • 一个NameNode + 若干DataNode
  • DataNode在NameNode统一调度下操作
  • 数据保存在本地Linux文件系统

3.4.2 命名空间管理

  • 命名空间包含:目录、文件和块
  • HDFS 1.0:整个集群只有一个命名空间,只有一个NameNode
  • 使用传统分级文件体系:支持创建、删除、转移、重命名

3.4.3 通信协议

  • 所有HDFS通信协议构建在TCP/IP基础之上

3.4.4 客户端

  • 用户与HDFS交互的接口

3.4.5 HDFS体系结构的局限性

  • HDFS只设置唯一一个NameNode,带来:
    • 单点故障风险
    • 内存受限(无法存储过多文件元数据)
    • 性能瓶颈(所有请求经过NameNode)

3.5 HDFS存储原理

3.5.1 冗余数据保存

  • 采用多副本方式冗余存储
  • 数据块的多个副本分布到不同DataNode
  • 优点:
    1. 加快数据传输速度(就近读取)
    2. 保证数据可用性和容错性
    3. 保证数据可靠性

3.5.2 数据存取策略

数据存放(副本放置策略)

  • 第一个副本:放在本地机架
  • 第二个副本:放在不同机架的节点
  • 第三个副本:放在与第二个副本相同机架的不同节点

数据读取

  • 优先读取距离客户端最近的副本

3.5.3 数据错误与恢复

名称节点出错

  • 备份机制:核心文件同步到SecondaryNameNode
  • 出错时根据备份恢复

数据节点出错

  • DataNode定期向NameNode发送**“心跳”信息**
  • 心跳中断 → NameNode标记该节点不可用
  • HDFS会自动在其他节点上重新创建丢失的副本

数据出错

  • 网络传输和磁盘错误可能造成数据损坏
  • 客户端读取时校验数据完整性(校验和CheckSum)
  • 发现损坏时从其他副本读取并修复

3.6 HDFS数据读写过程

读数据过程

  1. 客户端调用open()打开文件
  2. RPC请求NameNode获取文件块位置
  3. NameNode返回按距离排序的DataNode列表
  4. 客户端选择最近的DataNode读取数据
  5. 读取完毕后关闭连接

写数据过程

  1. 客户端调用create()创建文件
  2. RPC请求NameNode,NameNode检查文件是否存在
  3. NameNode返回可用的DataNode列表
  4. 客户端将数据分块写入DataNode
  5. DataNode之间建立Pipeline进行副本复制
  6. 完成后通知NameNode更新元数据

3.7 HDFS编程实践

HDFS常用命令

hadoop fs -ls <path>           # 显示文件详细信息
hadoop fs -ls -R <path>        # 递归显示
hadoop fs -cat <path>          # 输出文件内容
hadoop fs -mkdir [-p] <paths>  # 创建目录
hadoop fs -copyFromLocal <src> <dst>  # 上传文件
hadoop fs -copyToLocal <src> <localdst>  # 下载文件
hadoop fs -cp <src> <dst>       # 复制
hadoop fs -get <src> <localdst> # 下载
hadoop fs -rm <path>            # 删除文件
hadoop fs -rm -r <path>         # 递归删除

Java API编程

  • FileSystem.get(conf) → 获取文件系统实例
  • fs.open(new Path(uri)) → FSDataInputStream(读)
  • fs.create(new Path(uri)) → FSDataOutputStream(写)
  • fs.listStatus(path) → 列出文件状态
  • PathFilter 接口 → 自定义文件过滤

编程实例:文件过滤与合并

  • 过滤掉特定后缀文件(如.abc)
  • 合并剩余文件内容到一个文件
  • 涉及MyPathFilter、FSDataInputStream、FSDataOutputStream

第4章 分布式文件系统HDFS(下)

第4章与第3章内容高度重叠,为HDFS的补充和深化。重复知识点不再赘述,以下补充差异内容。

补充要点

HDFS 2.0的改进(联邦机制)

  • HDFS 1.0的单NameNode架构存在局限性
  • HDFS 2.0引入**联邦(Federation)**机制
  • 解决单NameNode的内存受限和性能瓶颈问题

HDFS高可用性(HA)

  • 引入Active/Standby双NameNode架构
  • 通过ZooKeeper实现自动故障转移

第5章 MapReduce

4.1 概述

4.1.1 分布式并行编程

  • 2005年后摩尔定律逐渐失效
  • 数据量快速增长,转向分布式并行编程提高性能
  • 谷歌最先提出MapReduce模型
  • Hadoop MapReduce是开源实现,使用门槛更低

4.1.2 MapReduce模型简介

  • 将复杂的并行计算高度抽象为两个函数:MapReduce
  • <k1,v1> → Map → <k2,v2> → Shuffle → <k2,list(v2)> → Reduce → <k3,v3>

4.1.3 Map和Reduce函数

  • Map:对输入的每个键值对应用一个函数,产生一组中间键值对
  • Reduce:对具有相同key的所有value进行归约操作
  • Combiner(可选):Map端的本地聚合(类似本地Reduce)

4.2 MapReduce体系结构

四个部分:

  1. Client:用户提交作业、查看运行状态
  2. JobTracker:管理所有作业,负责作业调度和任务监控
  3. TaskTracker:运行在从节点,执行具体任务,向JobTracker汇报
  4. Task:分为Map Task和Reduce Task,由TaskTracker启动

4.3 MapReduce工作流程

4.3.1 工作流程概述

  1. 提交作业
  2. JobTracker初始化作业
  3. 分配任务
  4. 执行Map任务
  5. Shuffle阶段
  6. 执行Reduce任务
  7. 作业完成

4.3.2 各执行阶段

Split(分片)

  • HDFS以block为存储单位,MapReduce以split为处理单位
  • split是逻辑概念,包含元数据(起始位置、长度、所在节点)
  • Hadoop为每个split创建一个Map任务

Map和Reduce任务数量

  • Map任务数 = split数(理想情况一个split = 一个HDFS块)
  • Reduce任务数 = 通常比可用slot数略少(预留资源处理错误)

4.3.3 Shuffle过程详解

Map端Shuffle

  1. Map输出写入环形缓冲区
  2. 缓冲区满后**溢写(Spill)**到磁盘
  3. 溢写前进行排序Combiner合并
  4. 多个溢写文件**归并(Merge)**成一个大文件

Combine和Merge的区别

  • Combine:<"a",1> + <"a",1><"a",2>(值相加)
  • Merge:<"a",1> + <"a",1><"a",<1,1>>(值放在一起)

Reduce端Shuffle

  1. Reduce任务从各Map任务拉取数据
  2. 对拉取的数据进行归并和排序
  3. 相同key的不同value按序归并,供Reduce使用

4.3.4 执行过程

  • 提交→初始化→分配→Map→Shuffle→Reduce→完成

4.4 实例分析:WordCount

设计思路

  • Map阶段:将每行文本拆分成单词,输出 <单词, 1>
  • Shuffle阶段:相同单词的value归并
  • Reduce阶段:对每个单词的所有1求和,输出 <单词, 总数>

执行过程

  • 有Combiner vs 无Combiner的对比:
    • 无Combiner:所有<word,1>都传到Reduce端
    • 有Combiner:Map端先做局部求和,减少网络传输

4.5 MapReduce的具体应用

应用类型

  • 排序
  • 数据挖掘
  • 日志分析
  • 关系自然连接

用MapReduce实现关系自然连接

  • 设计思路:利用Map端的key分组,相同连接属性的元组会被分配到同一个Reduce任务
  • Reduce端完成实际的连接操作

4.6 MapReduce编程实践

编写Map处理逻辑

public static class MyMapper extends Mapper<Object,Text,Text,IntWritable>{
    private final static IntWritable one = new IntWritable(1);
    private Text word = new Text();
    public void map(Object key, Text value, Context context)
        throws IOException,InterruptedException{
        StringTokenizer itr = new StringTokenizer(value.toString());
        while (itr.hasMoreTokens()){
            word.set(itr.nextToken());
            context.write(word,one);
        }
    }
}

编写Reduce处理逻辑

public static class MyReducer extends Reducer<Text,IntWritable,Text,IntWritable>{
    private IntWritable result = new IntWritable();
    public void reduce(Text key, Iterable<IntWritable> values, Context context)
        throws IOException,InterruptedException{
        int sum = 0;
        for (IntWritable val : values){
            sum += val.get();
        }
        result.set(sum);
        context.write(key,result);
    }
}

编写main方法

public static void main(String[] args) throws Exception{
    Configuration conf = new Configuration();
    String[] otherArgs = new GenericOptionsParser(conf,args).getRemainingArgs();
    Job job = new Job(conf,"word count");
    job.setJarByClass(WordCount.class);
    job.setMapperClass(MyMapper.class);
    job.setReducerClass(MyReducer.class);
    job.setOutputKeyClass(Text.class);
    job.setOutputValueClass(IntWritable.class);
    FileInputFormat.addInputPath(job,new Path(otherArgs[0]));
    FileOutputFormat.setOutputPath(job,new Path(otherArgs[1]));
    System.exit(job.waitForCompletion(true)?0:1);
}

编译打包运行

  • 设置CLASSPATH:export CLASSPATH=$($HADOOP_HOME/bin/hadoop classpath):$CLASSPATH
  • 编译:javac WordCount.java
  • 打包:jar cvf WordCount.jar *.class
  • 运行:hadoop jar WordCount.jar input output

第6章 Spark

6.1 Spark概述

6.1.1 Spark简介

  • 2009年由加州伯克利大学AMP实验室开发
  • 基于内存计算的大数据并行计算框架
  • 2014年打破Hadoop的基准排序纪录

Spark特点

  • 运行速度快(内存计算)
  • 易于使用
  • 支持通用计算模型
  • 运行模式多样

6.1.2 Scala简介

  • 现代多范式编程语言,运行于JVM
  • Spark的主要编程语言(也支持Java、Python、R)
  • 提供REPL(Read-Eval-Print Loop),提高开发效率

6.1.3 Spark与Hadoop对比

Hadoop缺点

  • 表达能力有限(只有Map和Reduce两个操作)
  • 磁盘I/O开销大(中间结果写入磁盘)
  • 执行延迟高(每个Job都要读写磁盘)

Spark优点

  • 基于内存计算,中间结果保存在内存中
  • 丰富的数据操作(转换+行动)
  • 更快的执行速度(尤其是迭代计算)

Hadoop vs Spark 迭代计算对比

  • Hadoop:每次迭代都要读写磁盘,非常耗资源
  • Spark:数据载入内存后,后续迭代直接使用内存中间结果

6.2 Spark生态系统

组件组成

组件功能
Spark Core核心引擎,RDD、任务调度
Spark SQL结构化数据处理,SQL查询
Spark Streaming流式数据处理(微批)
Structured Streaming持续流式处理(毫秒级)
MLlib机器学习库
GraphX图计算

BDAS(Berkeley Data Analytics Stack)

  • Spark生态系统是BDAS的重要组成部分

应用场景

  • 批处理、交互式查询、流数据处理、机器学习、图计算

6.3 Spark运行架构

6.3.1 基本概念

  1. Application:用户编写的Spark应用程序
  2. Driver:运行Application的main()函数,创建SparkContext
  3. Cluster Manager:集群资源管理器(Standalone、YARN、Mesos、K8s)
  4. Executor:工作节点上为Application启动的进程,执行Task
  5. Worker Node:集群中运行Application代码的节点
  6. Task:被分配到Executor上的最小工作单元
  7. Job:由多个Stage组成的并行计算
  8. Stage:由多个没有Shuffle关系的Task组成
  9. RDD:弹性分布式数据集

6.3.2 架构设计

  • Driver + Cluster Manager + Executor
  • Executor的两个优点:多线程执行Task、数据本地化优化
  • Application = Driver + 多个Job
  • Job = 多个Stage
  • Stage = 多个Task(无Shuffle关系)

6.3.3 Spark运行基本流程

  1. 构建SparkContext(连接Cluster Manager)
  2. Cluster Manager分配资源,启动Executor
  3. Driver发送代码和文件到Executor
  4. Executor执行Task,结果返回Driver或写入外部存储

6.3.4 RDD的设计与运行原理

RDD概念

  • 弹性分布式数据集(Resilient Distributed Dataset)
  • 只读的记录分区的集合
  • 创建方式:
    1. 从稳定物理存储创建(如HDFS)
    2. 从其他RDD通过转换操作创建(map、join、group by等)

RDD执行过程

  • 惰性调用:转换操作只记录轨迹,不触发计算
  • 行动操作才触发真正的计算
  • 记录RDD之间的依赖关系 → 形成血缘关系(Lineage)
  • 根据依赖关系生成DAG(有向无环图)

RDD特性(实现高效计算的原因):

  • 高效容错:通过血缘关系恢复丢失分区(不需要数据复制)
  • 弹性:RDD是只读的,但如果数据丢失,可通过Lineage重新计算
  • 分区:数据被分成多个分区,可并行处理

RDD之间的依赖关系

Shuffle操作

  • Shuffle过程产生大量网络传输磁盘I/O开销
  • Map端:多个桶写入一个文件(优化磁盘I/O)
  • Reduce端:拉取数据、归并排序
  • 建议:涉及Shuffle操作时增加分区数量

窄依赖与宽依赖

类型定义示例操作
窄依赖父RDD的每个分区最多被子RDD的一个分区使用map、filter、flatMap
宽依赖父RDD的每个分区被子RDD的多个分区使用groupByKey、reduceByKey、join

阶段的划分(Stage)

  • Spark根据DAG中的RDD依赖关系划分Stage
  • 宽依赖处切分Stage(产生Shuffle)
  • 窄依赖实现流水线优化(Pipeline)
  • 宽依赖无法实现流水线优化

Stage划分规则:从后往前,遇到宽依赖就切一个新的Stage

6.4 Spark的部署和应用方式

部署方式

模式说明
Local单机模式
StandaloneSpark自带集群管理
Spark on YARN运行在YARN上
Spark on Mesos运行在Mesos上
Spark on K8s运行在Kubernetes上(2.3.0+)

从”Hadoop+Storm”转向Spark

  • 传统Lambda架构:Hadoop负责批处理 + Storm负责流处理
  • Spark同时支持批处理和流处理,简化部署
  • 注意:Spark Streaming本质是微批处理,非真正实时(无法毫秒级响应)
  • 毫秒级需求仍需Storm/Flink

Hadoop和Spark统一部署

  • 多个框架(MapReduce、Spark、HBase、Storm等)统一运行在YARN上
  • 好处:统一资源管理、共享底层存储、降低运维成本

6.5 Spark编程实践

Spark安装

tar -zxf spark-3.4.0-bin-without-hadoop.tgz -C /usr/local/
mv spark-3.4.0-bin-without-hadoop spark
# 配置Classpath
export SPARK_DIST_CLASSPATH=$(/usr/local/hadoop/bin/hadoop classpath)

RDD基本操作

转换操作(Transformation,惰性求值)

操作说明
map(func)每个元素应用func
flatMap(func)每个元素映射到0或多个输出
filter(func)筛选满足条件的元素
groupByKey()按key分组
reduceByKey(func)按key聚合

行动操作(Action,触发计算)

操作说明
count()返回元素个数
first()返回第一个元素
take(n)返回前n个元素
reduce(func)聚合所有元素
collect()返回所有元素到Driver
foreach(func)对每个元素应用func

用sbt编译打包Scala程序

  • 安装sbt
  • 创建项目结构:~/sparkapp/src/main/scala/
  • 编写Scala代码
  • 创建simple.sbt声明依赖
  • 打包:sbt package
  • 提交:spark-submit --class "SimpleApp" xxx.jar

用Maven编译打包Java程序

  • 安装Maven
  • 创建项目结构
  • 编写Java代码
  • 创建pom.xml声明依赖
  • 打包:mvn package
  • 提交:spark-submit --class "SimpleApp" xxx.jar

第7章 流计算

7.1 流计算概述

7.1.1 静态数据和流数据

静态数据

  • 历史数据,存储在数据仓库中
  • 可用数据挖掘和OLAP分析

流数据

  • 以大量、快速、时变的流形式持续到达
  • 例:PM2.5检测、电子商务用户点击流
  • 特征:大量、快速、时变、潜在价值巨大

7.1.2 批量计算和实时计算

  • 批量计算:处理静态数据(Hadoop MapReduce)
  • 实时计算:处理流数据(流计算框架)

7.1.3 流计算概念

  • 定义:实时获取来自不同数据源的海量数据,经实时分析处理,获得有价值信息
  • 基本理念:数据价值随时间流逝而降低,应立即处理
  • 系统需求:低延迟、可扩展、高可靠

7.1.4 流计算与Hadoop

  • Hadoop擅长批处理,不适合流计算
  • 两者互补,各有适用场景

7.1.5 流计算框架

三类:

  1. 商业级流计算平台
  2. 开源流计算框架(Storm、Flink、Spark Streaming)
  3. 公司自研流计算框架

7.2 流计算处理流程

7.2.1 数据处理流程

传统流程:采集→存储→查询(两个前提:数据有限、响应及时)

流计算流程(三个阶段):

  1. 数据实时采集
  2. 数据实时计算
  3. 实时查询服务

7.2.2 数据实时采集

  • 采集多个数据源的海量数据
  • 保证实时性、低延迟、稳定可靠
  • 开源系统:Flume、Kafka、Logstash等
  • 数据采集系统基本架构:Agent → Collector → Storage

7.2.3 数据实时计算

  • 对采集的数据进行实时分析和计算
  • 结果可视情况存储或直接丢弃
  • 高时效场景:处理后数据直接丢弃

7.2.4 实时查询服务

  • 传统:用户主动查询
  • 流处理:实时推送结果给用户
  • 本质区别:流处理结果是实时的,传统定时查询结果是过去的

7.3 流计算的应用

应用场景

  1. 实时分析:如量子恒道Super Mario框架,TB级/天实时流处理,2-3秒延迟
  2. 实时交通:根据实时交通数据更新导航路线

7.4 Storm

  • 2011年Twitter开源的流计算框架
  • 支持真正毫秒级实时响应

7.5 Spark Streaming

原理

  • 将实时数据流按时间片(秒级)拆分
  • 每个时间片作为一个批处理作业处理
  • 微批处理,非真正实时

核心抽象

  • DStream(Discretized Stream,离散化数据流)
  • 输入数据按时间片分成一段一段的DStream
  • 每段数据转换为RDD
  • 对DStream的操作最终转变为对RDD的操作

数据源

  • 支持Kafka、Flume、HDFS、TCP套接字等多种输入源
  • 输出可存储至文件系统、数据库或显示在仪表盘

7.6 Structured Streaming

7.6.1 简介

  • Spark 2.3.0引入的持续流式处理模型
  • 可将延迟降低到毫秒级别

7.6.2 关键思想

  • 将实时数据流视为一张正在不断添加数据的表(无界表)
  • 流计算 = 在不断添加数据的无界输入表上运行增量查询
  • 输入流的每项数据被添加到无界表,形成新的无界表

7.6.3 两种处理模型

  1. 微批处理(默认)

    • 将数据分成小批次处理
    • 延迟:通常百毫秒级别
    • 适合大多数场景(ETL、监控等)
  2. 持续处理

    • 真正的逐条处理
    • 延迟:10~20ms
    • 适合低延迟场景(如信用卡欺诈检测)
    • 从Spark 2.3.0开始支持
  • 真正的流计算框架,原生支持事件时间处理
  • 支持有状态计算,精确一次语义
  • 低延迟、高吞吐、高可靠

第8章 大数据应用

8.1 互联网领域

推荐系统

  • 本质:建立用户与物品的联系
  • 推荐方法:基于内容、协同过滤、混合推荐等

8.2 生物医学领域

8.2.1 流行病预测

  • 谷歌流感趋势:通过搜索词数据判断流感情况

8.2.2 智慧医疗

  • 促进优质医疗资源共享
  • 促进医疗智能化
  • 避免患者重复检查

8.2.3 生物信息学

  • 研究生物信息的采集、处理、存储、传播、分析和解释
  • 生命科学与计算机科学的交叉学科
  • 生物大数据应用:了解生物学过程、作物表型、疾病致病基因
  • 从个人健康档案预测健康趋势,“治未病”

8.3 物流领域

8.3.1 智能物流

  • 利用集成智能化技术,模仿人的智能
  • 实现物流资源优化调度和有效配置

8.3.2 大数据是智能物流的关键

  • 两个著名理论:“黑大陆说””物流冰山说”
  • 大数据推动物流从粗放式到个性化服务转变

8.3.3 菜鸟网络

  • 中国智能物流骨干网
  • 目标:24小时内送达国内任何地区
  • 支撑日均300亿元网络零售额
  • 1000亿元投资物流基础设施

8.4 城市管理领域

8.4.1 智能交通

  • 集成信息技术、数据通信、电子传感、控制、计算机技术
  • 利用实时交通信息、社交网络和天气数据优化交通

8.4.2 环保监测

  • 森林监视:谷歌森林监视
  • 环境保护:污染监测、中国水/空气/固废污染地图、汽车尾气治理

8.4.3 城市规划

  • 大数据辅助城市规划决策

8.4.4 安防

  • 平安城市建设:密布摄像头7×24小时采集视频数据
  • 数据量大,包含结构化、半结构化和非结构化数据

8.4.5 疫情防控

  • 2020年新冠疫情:大数据在疫情防控、资源调配、复工复产中发挥重要作用

8.5 金融领域

8.5.1 高频交易(HFT)

  • 从极为短暂的市场变化中寻求获利
  • 引入大数据技术辅助交易决策

8.5.2 市场情绪分析

  • 大数据技术在市场情绪分析中大有用武之地

8.5.3 信贷风险分析

  • 收集分析中小微企业日常交易数据
  • 判断业务、经营、信用状况
  • 解决财务制度不健全导致的评估难题

8.5.4 大数据征信

  • 将不同信贷机构、消费场景的海量数据整合
  • 经过数据清洗、模型分析、校验加工成有用信息

8.6 汽车领域

  • 无人驾驶汽车:大量传感器每秒产生1GB数据,每年约2PB
  • 大数据帮助做出更智能的驾驶决策
  • 比人类驾车更安全、舒适、节能、环保

8.7 零售领域

8.7.1 发现关联购买行为

  • 啤酒与尿布:沃尔玛经典案例,关联规则挖掘

8.7.2 客户群体细分

  • Target超市比孩子父亲更早发现女儿怀孕(数据挖掘案例)

8.7.3 供应链管理

  • 亚马逊、UPS、沃尔玛利用大数据掌控供应链
  • McKesson:每天200万订单全程跟踪分析,监督80亿美元存货

8.8 餐饮领域

  • Food Genius:聚合35万家餐馆菜单数据,帮助确定价格、趋势
  • 餐饮O2O:大数据驱动团购、推荐、门店布局、人流量控制

8.9 电信领域

  • 电信客户离网分析:预测客户行为、发现行为趋势、找出服务缺陷
  • 及时采取措施保留客户

8.10 能源领域

  • 智能电网离不开大数据技术
  • 大数据是智能电网的技术基石

8.11 体育和娱乐领域

8.11.1 训练球队

  • 大数据帮助球队提升实力和水平

8.11.2 投拍影视作品

  • 大数据分析辅助影视投资决策

8.11.3 预测比赛结果

  • 基于大数据和预测模型预测事件概率
  • 对人们生活产生深刻影响

8.12 安全领域

8.12.1 国家安全

  • 大数据在国家安全中的应用

8.12.2 网络攻击防御

  • 云杀毒:客户端监测→云端分析→解决方案分发
  • 利用云计算和大数据及时发现未知威胁

8.12.3 预防犯罪

  • 洛杉矶警察局:首个大数据公安警务模式
  • 根据大数据预测犯罪、提供破案线索、实时犯罪预警

8.13 日常生活

  • 个人大数据:通话、聊天、邮件、购物、出行、住宿、生理指标
  • 分析个人大数据了解生活行为模式
  • 为个人提供更加周到的服务

复习文档完成 ✧ 覆盖全部8章内容,包含:大数据概述(4V特征、三种思维方式变革)、数据存储与管理(关系数据库→NoSQL→NewSQL→数据湖→湖仓一体)、HDFS(块、NameNode/DataNode、FsImage/EditLog、数据读写、编程实践)、MapReduce(Map/Reduce/Shuffle/WordCount)、Spark(RDD/窄宽依赖/Stage划分/DAG/部署方式/编程)、流计算(Spark Streaming/Structured Streaming/Flink)、大数据应用(12个领域)。