大数据技术 系统性复习文档
系统性复习文档:涵盖大数据概述、数据存储与管理、分布式文件系统 HDFS、MapReduce 等核心章节。
大数据技术 系统性复习文档
目录
- 第1章 大数据概述
- 第2章 数据存储与管理
- 第3章 分布式文件系统HDFS(上)
- 第4章 分布式文件系统HDFS(下)
- 第5章 MapReduce
- 第6章 Spark
- 第7章 流计算
- 第8章 大数据应用
第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 数据产生方式的变革
三个阶段:
- 运营式系统阶段:数据被动产生,伴随运营活动存入数据库
- 用户原创内容阶段:数据主动产生(Web 2.0时代),智能手机加速内容产生
- 感知式系统阶段:传感器广泛使用,数据自动采集,促成大数据产生
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 大数据的应用
应用三个层次
- 描述性分析:总结抽取信息,分析”发生了什么”
- 例:DOMO公司从企业数据中提取信息推送
- 预测性分析:分析关联关系和发展模式,预测趋势
- 例:David Rothschild用大数据预测奥斯卡(21/24准确)
- 指导性分析:分析决策后果,指导优化决策
- 例:无人驾驶汽车实时决策
典型应用实例
- 《纸牌屋》:大数据分析决定选角和内容
- 谷歌流感趋势:跟踪搜索词判断流感情况
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 数据的概念
文件系统
- 操作系统中负责管理和存储文件信息的软件机构
- 由三部分组成:文件系统接口、软件集合、对象及属性
- 负责为用户建立、存入、读出、修改、转储、撤销文件
数据库
- 按一定规则组织、可共享、冗余少、与程序分开的数据集合
- 特点:按一定方式存储、多用户共享、冗余度小、与应用程序独立
数据库发展历史
- 层次数据库(树形结构):简单但关系不够灵活
- 网状数据库:多对多联系,复杂难用
- 关系数据库:二维表存储,简单清晰易理解
关系数据库
- 用二维表(行和列)组织数据
- 例:学生信息表
数据仓库
- 定义:面向主题的、集成的、相对稳定的、反映历史变化的数据集合,用于支持管理决策
- 数据仓库体系架构
并行数据库
- 在无共享体系结构中进行数据操作的数据库系统
- 采用关系数据模型,支持SQL
- 两个关键技术:关系表的水平划分、SQL查询的分区执行
- 目标:高性能和高可用性
- 缺点:数据转移代价昂贵
2.2 大数据时代数据存储与管理技术
分布式文件系统
- 通过网络实现文件在多台主机上分布式存储
- GFS:谷歌开发的分布式文件系统
- HDFS:GFS的开源实现,Hadoop两大核心之一
NewSQL和NoSQL数据库
NewSQL数据库:
- 新的可扩展、高性能数据库
- 兼具海量数据存储能力和传统数据库的ACID、SQL支持
- 代表产品:Spanner、VoltDB、Clustrix等
NoSQL数据库:
- 非关系型数据库的统称
- 没有固定表结构,无连接操作,不严格遵守ACID
- 灵活的水平可扩展性,支持海量数据
- 数据模型:键/值、列族、文档等非关系模型
大数据引发的数据库架构变革:
- 从单机关系数据库 → 并行数据库 → 分布式数据库 → NoSQL数据库
云数据库
- 云服务提供商在公有云中托管数据库
- 满足企业动态数据存储需求和中小企业低成本存储需求
- 数据量每年60%速度增加
数据湖
- 定义:企业中全量数据的单一存储
- 包含:结构化、半结构化、非结构化、二进制数据
- 可构建在本地数据中心或云上
- 本质:数据存储架构 + 数据处理工具
- 存储底座:通常用对象存储(如Amazon S3)
- 两类工具:
- 数据移动和管理工具(搬到湖里)
- 数据分析工具(从湖里淘金)
数据湖与数据仓库的区别:
- 数据仓库:结构化数据,模式固定,适合报表和分析
- 数据湖:全量数据,模式灵活,适合探索和机器学习
湖仓一体:
- 打通数据仓库和数据湖的新型开放式架构
- 融合数据仓库的高性能管理能力与数据湖的灵活性
- 特性:存算分离、事务支持、开放性、数据治理、支持多种数据类型、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)
- 文件被分成多个块,以块为存储单位
- 块的优点:
- 支持大规模文件存储
- 简化系统设计(固定大小简化元数据管理)
- 适合数据备份(利于容错和数据可用性)
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
- 优点:
- 加快数据传输速度(就近读取)
- 保证数据可用性和容错性
- 保证数据可靠性
3.5.2 数据存取策略
数据存放(副本放置策略):
- 第一个副本:放在本地机架
- 第二个副本:放在不同机架的节点
- 第三个副本:放在与第二个副本相同机架的不同节点
数据读取:
- 优先读取距离客户端最近的副本
3.5.3 数据错误与恢复
名称节点出错:
- 备份机制:核心文件同步到SecondaryNameNode
- 出错时根据备份恢复
数据节点出错:
- DataNode定期向NameNode发送**“心跳”信息**
- 心跳中断 → NameNode标记该节点不可用
- HDFS会自动在其他节点上重新创建丢失的副本
数据出错:
- 网络传输和磁盘错误可能造成数据损坏
- 客户端读取时校验数据完整性(校验和CheckSum)
- 发现损坏时从其他副本读取并修复
3.6 HDFS数据读写过程
读数据过程
- 客户端调用
open()打开文件 - RPC请求NameNode获取文件块位置
- NameNode返回按距离排序的DataNode列表
- 客户端选择最近的DataNode读取数据
- 读取完毕后关闭连接
写数据过程
- 客户端调用
create()创建文件 - RPC请求NameNode,NameNode检查文件是否存在
- NameNode返回可用的DataNode列表
- 客户端将数据分块写入DataNode
- DataNode之间建立Pipeline进行副本复制
- 完成后通知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模型简介
- 将复杂的并行计算高度抽象为两个函数:Map 和 Reduce
<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体系结构
四个部分:
- Client:用户提交作业、查看运行状态
- JobTracker:管理所有作业,负责作业调度和任务监控
- TaskTracker:运行在从节点,执行具体任务,向JobTracker汇报
- Task:分为Map Task和Reduce Task,由TaskTracker启动
4.3 MapReduce工作流程
4.3.1 工作流程概述
- 提交作业
- JobTracker初始化作业
- 分配任务
- 执行Map任务
- Shuffle阶段
- 执行Reduce任务
- 作业完成
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:
- Map输出写入环形缓冲区
- 缓冲区满后**溢写(Spill)**到磁盘
- 溢写前进行排序和Combiner合并
- 多个溢写文件**归并(Merge)**成一个大文件
Combine和Merge的区别:
- Combine:
<"a",1>+<"a",1>→<"a",2>(值相加)- Merge:
<"a",1>+<"a",1>→<"a",<1,1>>(值放在一起)
Reduce端Shuffle:
- Reduce任务从各Map任务拉取数据
- 对拉取的数据进行归并和排序
- 相同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端先做局部求和,减少网络传输
- 无Combiner:所有
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 基本概念
- Application:用户编写的Spark应用程序
- Driver:运行Application的main()函数,创建SparkContext
- Cluster Manager:集群资源管理器(Standalone、YARN、Mesos、K8s)
- Executor:工作节点上为Application启动的进程,执行Task
- Worker Node:集群中运行Application代码的节点
- Task:被分配到Executor上的最小工作单元
- Job:由多个Stage组成的并行计算
- Stage:由多个没有Shuffle关系的Task组成
- RDD:弹性分布式数据集
6.3.2 架构设计
- Driver + Cluster Manager + Executor
- Executor的两个优点:多线程执行Task、数据本地化优化
- Application = Driver + 多个Job
- Job = 多个Stage
- Stage = 多个Task(无Shuffle关系)
6.3.3 Spark运行基本流程
- 构建SparkContext(连接Cluster Manager)
- Cluster Manager分配资源,启动Executor
- Driver发送代码和文件到Executor
- Executor执行Task,结果返回Driver或写入外部存储
6.3.4 RDD的设计与运行原理
RDD概念:
- 弹性分布式数据集(Resilient Distributed Dataset)
- 只读的记录分区的集合
- 创建方式:
- 从稳定物理存储创建(如HDFS)
- 从其他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 | 单机模式 |
| Standalone | Spark自带集群管理 |
| 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 流计算框架
三类:
- 商业级流计算平台
- 开源流计算框架(Storm、Flink、Spark Streaming)
- 公司自研流计算框架
7.2 流计算处理流程
7.2.1 数据处理流程
传统流程:采集→存储→查询(两个前提:数据有限、响应及时)
流计算流程(三个阶段):
- 数据实时采集
- 数据实时计算
- 实时查询服务
7.2.2 数据实时采集
- 采集多个数据源的海量数据
- 保证实时性、低延迟、稳定可靠
- 开源系统:Flume、Kafka、Logstash等
- 数据采集系统基本架构:Agent → Collector → Storage
7.2.3 数据实时计算
- 对采集的数据进行实时分析和计算
- 结果可视情况存储或直接丢弃
- 高时效场景:处理后数据直接丢弃
7.2.4 实时查询服务
- 传统:用户主动查询
- 流处理:实时推送结果给用户
- 本质区别:流处理结果是实时的,传统定时查询结果是过去的
7.3 流计算的应用
应用场景
- 实时分析:如量子恒道Super Mario框架,TB级/天实时流处理,2-3秒延迟
- 实时交通:根据实时交通数据更新导航路线
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 两种处理模型
-
微批处理(默认):
- 将数据分成小批次处理
- 延迟:通常百毫秒级别
- 适合大多数场景(ETL、监控等)
-
持续处理:
- 真正的逐条处理
- 延迟:10~20ms
- 适合低延迟场景(如信用卡欺诈检测)
- 从Spark 2.3.0开始支持
7.7 Flink
- 真正的流计算框架,原生支持事件时间处理
- 支持有状态计算,精确一次语义
- 低延迟、高吞吐、高可靠
第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个领域)。