1.Spark 核心架构 1.1 整体架构 ┌─────────────────────────────────────────────────────────
│ Client(客户端)
│ spark-submit / pyspark / spark-shell
└─────────────────────────────────────────────────────────
│
▼
┌─────────────────────────────────────────────────────────
│ Driver(驱动器)
│ ┌─────────────┐ ┌──────────────┐ ┌───────────────┐
│ │ DAGScheduler │ │ Scheduler │ │ SparkContext
│ │ (任务调度) │ │ (资源调度) │ │ (上下文管理) │
│ └─────────────┘ └──────────────┘ └───────────────┘
└─────────────────────────────────────────────────────────
│
┌───────────────┼───────────────┐
▼ ▼ ▼
┌───────────┐ ┌───────────┐ ┌───────────┐
│ Executor │ │ Executor │ │ Executor │ ← Worker节点
│ (工作节点) │ │ (工作节点) │ │ (工作节点) │
└───────────┘ └───────────┘ └───────────┘
│ │ │
└───────────────┴───────────────┘
│
┌───────────┴───────────┐
▼ ▼
┌─────────┐ ┌─────────┐
│ HDFS │ │ Kafka │ ← 数据源
└─────────┘ └─────────┘
1.2核心组件详解
组件
职责
说明
SparkContext
集群连接的入口
创建RDD、广播变量、累加器
DAGScheduler
阶段调度器
将Job划分为Stage,提交到TaskScheduler
TaskScheduler
任务调度器
将Task分配到具体Executor执行
Cluster Manager
资源管理器
支持Standalone/YARN/Mesos/Kubernetes
Executor
执行器
运行Task、缓存数据、提供内存池
2.计算模型:RDD RDD(弹性分布式数据集) 1 2 3 4 rdd = sc.parallelize([1 ,2 ,3 ,4 ,5 ]) rdd2 = rdd.map (lambda x: x*2 ) result = rdd2.filter (lambda x: x>5 ).collect()
RDD特性
特性
说明
弹性
内存与磁盘自动交换、失败自动重试
分布式
数据分片分布在多个节点
只读
不可变,通过转换生成新RDD
分区
数据按分区并行处理
RDD依赖关系 Stage划分示例:
BroadCasta → [D] → [E] → [F]
↑ ↑
[A] → [B] → [C] ← Shuffle依赖 → 新Stage
窄依赖 :父RDD一个分区最多被子RDD一个分区使用
宽依赖 :父RDD一个分区需被子RDD多个分区使用(触发Shuffle)
3.核心模块 3.1Spark Core
RDD基础抽象
内存管理、调度、故障恢复
任务调度与执行引擎
3.2 Spark SQL 1 2 3 4 5 df = spark.read.json("hdfs://path/to/data") df.createOrReplaceTempView("people") result = spark.sql("SELECT name, age FROM people WHERE age > 30")result.write.parquet("hdfs://path/to/output")
对比项
DataFrame
Dataset
类型安全
❌
✅
性能
高(Catalyst优化器)
更高(编译时检查)
语言支持
多语言
JVM语言为主
3.3 Spark Streaming 1 2 3 4 5 6 7 8 ssc = StreamingContext(spark.sparkContext, batchDuration=10 ) lines = ssc.socketTextStream("localhost" , 9999 ) words = lines.flatMap(lambda line: line.split(" " )) wordCounts = words.map (lambda x: (x, 1 )).reduceByKey(lambda a,b: a+b) wordCounts.pprint() ssc.start() ssc.awaitTermination()
3.4 MLlib(机器学习) 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 from pyspark.ml.classification import LogisticRegressionfrom pyspark.ml.feature import VectorAssemblerassembler = VectorAssembler(inputCols=["age" , "income" ], outputCol="features" ) data = assembler.transform(df) lr = LogisticRegression(maxIter=10 , regParam=0.01 ) model = lr.fit(data) predictions = model.transform(data) evaluator = MulticlassClassificationEvaluator() accuracy = evaluator.evaluate(predictions)
3.5 GraphX(图计算)
场景
示例
社交网络
好友推荐、影响力分析
知识图谱
实体关系推理
推荐系统
协同过滤
4.执行模式 4.1 集群管理器对比
模式
部署复杂度
资源隔离
适用场景
Standalone
低
中等
小规模集群、测试环境
YARN
中
强
Hadoop生态、企业生产
Mesos
高
强
异构集群、多框架混部
Kubernetes
中高
强
云原生、容器化部署
4.2 部署命令 1 2 3 4 5 6 7 8 9 10 11 12 13 14 spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 10 \ --executor-memory 4g \ --class com.example.Main app.jar spark-submit \ --master k8s://https://k8s-api:6443 \ --deploy-mode cluster \ --conf spark.kubernetes.container.image=spark:latest \ app.jar
5.性能优化 5.1 配置调优
参数
推荐值
说明
spark.executor.memory
4-8g
Executor堆内存
spark.executor.cores
2-4
每个Executor的CPU核数
spark.sql.shuffle.partitions
200-500
Shuffle分区数
spark.default.parallelism
等于总CPU核数
默认并行度
5.2 数据倾斜处理 1 2 3 4 5 6 7 8 9 10 11 from pyspark.sql.functions import concat, lit, monotonically_increasing_idsalted_df = df.withColumn("salt" , monotonically_increasing_id() % 10 ) result = salted_df.groupby("key" , "salt" ).agg(...) spark.sparkContext.broadcast(largeRDD) first_agg = df.groupby("key" ).agg(sum (...).alias("sum_v" )) second_agg = df.join(first_agg, "key" ).agg(sum (col("v" )*col("sum_v" )))
5.3 存储优化
格式
压缩
适用场景
Parquet
Snappy
列式存储,查询性能好
ORC
ZLIB
Hive生态、复杂查询
Avro
Snappy
行式存储,写多读少
6.应用场景 6.1 离线批处理(Spark Batch) 典型架构:HDFS → Spark → Hive/ClickHouse
场景
说明
数据仓库ETL
每日T+1报表、指标计算
用户画像构建
标签体系、人群圈选
日志分析
访问量统计、异常检测
6.2 实时计算(Spark Streaming) 典型架构:Kafka → Spark Streaming → Redis/ES
场景
说明
实时指标监控
DAU/UV、订单量统计
实时推荐
基于用户行为的实时推荐
实时告警
异常交易、系统故障告警
6.3 机器学习(MLlib)
场景
算法
信用评分
逻辑回归、GBDT
推荐系统
ALS、协同过滤
分类预测
决策树、随机森林
聚类分析
K-Means、Gaussian Mixture
6.4 图计算(GraphX)
场景
应用
社交网络
社区发现、影响力分析
知识图谱
路径查找、关系推理
风控网络
团伙识别、关联分析
7.与其他计算框架对比
框架
计算模型
延迟
吞吐量
适用场景
Spark
RDD/DAG
秒级
高
离线批处理、交互式查询
Flink
DataStream
毫秒级
中
实时流处理、复杂事件处理
Storm
Topology
毫秒级
低
简单实时计算
MapReduce
Map/Reduce
分钟级
极高
超大规模离线批处理