Spark技术架构及应用场景

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基本操作示例
rdd = sc.parallelize([1,2,3,4,5]) # 创建RDD
rdd2 = rdd.map(lambda x: x*2) # 转换操作(Transformation)
result = rdd2.filter(lambda x: x>5).collect() # 行动操作(Action)

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
-- DataFrame/Dataset API示例
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 LogisticRegression
from pyspark.ml.feature import VectorAssembler

# 特征组装
assembler = 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 on YARN
spark-submit \
--master yarn \
--deploy-mode cluster \
--num-executors 10 \
--executor-memory 4g \
--class com.example.Main app.jar

# Spark on Kubernetes
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
# 方案一:加盐(Salting)
from pyspark.sql.functions import concat, lit, monotonically_increasing_id
salted_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 分钟级 极高 超大规模离线批处理