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内存配置

更多推荐