精通Spark算子:性能优化实战指南,【Linux】操作系统的认识。
·
Spark算子介绍
Spark算子分为Transformation(转换)和Action(执行)两大类,用于处理分布式数据集的RDD操作。转换算子执行惰性计算,仅记录操作不立即执行;执行算子触发实际计算并返回结果。
常见Transformation算子
map:对RDD每个元素应用函数,返回新RDD
val rdd = sc.parallelize(Array(1, 2, 3))
val mapped = rdd.map(x => x * 2) // 2,4,6
filter:筛选符合条件的元素
val filtered = rdd.filter(x => x > 1) // 2,3
flatMap:先映射后扁平化
val fm = rdd.flatMap(x => Array(x, x*10)) // 1,10,2,20,3,30
groupByKey:按key分组
val kv = sc.parallelize(Array(("a",1),("b",2),("a",3)))
kv.groupByKey() // ("a",(1,3)), ("b",(2))
reduceByKey:按键聚合
kv.reduceByKey(_ + _) // ("a",4), ("b",2)
常见Action算子
collect:返回所有元素到driver
rdd.collect() // Array(1,2,3)
count:统计元素数量
rdd.count() // 3
reduce:聚合所有元素
rdd.reduce(_ + _) // 6
saveAsTextFile:保存到文件系统
rdd.saveAsTextFile("hdfs://path")
foreach:对每个元素执行操作
rdd.foreach(println)
关键特性对比
转换算子特点:
- 惰性执行机制
- 可形成血缘关系(lineage)
- 支持checkpoint容错
执行算子特点:
- 触发Job提交
- 返回具体数值或外部存储
- 会切断RDD依赖链
优化建议
选择reduceByKey而非groupByKey:
- reduceByKey会在map端合并
- 减少shuffle数据量
合理设置并行度:
- partition数量影响性能
- 可通过repartition调整
避免数据倾斜:
- 使用salting技术
- 考虑二次聚合
缓存复用RDD:
- persist()缓存中间结果
- 选择合适的存储级别
高级算子
aggregateByKey:更灵活的聚合
val data = sc.parallelize(List(("a",3),("a",2),("b",4)))
data.aggregateByKey(0)(_+_, _+_) // ("a",5),("b",4)
combineByKey:自定义聚合逻辑
val res = data.combineByKey(
(v) => (v,1),
(acc:(Int,Int),v) => (acc._1+v,acc._2+1),
(acc1:(Int,Int),acc2:(Int,Int)) => (acc1._1+acc2._1, acc1._2+acc2._2)
)
注意事项
窄依赖与宽依赖:
- map/filter产生窄依赖
- groupByKey产生宽依赖
任务调度优化:
- 合理设置并行度
- 监控任务执行情况
内存管理:
- 避免OOM异常
- 调整executor内存配置
更多推荐
所有评论(0)