Spark算子简介

Spark算子是Spark框架中用于数据转换和处理的核心操作,主要分为两类:转换算子(Transformations)和行动算子(Actions)。转换算子用于数据集的转换,生成新的RDD;行动算子触发实际计算并返回结果。理解这些算子是高效使用Spark的基础。

转换算子(Transformations)

map(func) 对RDD中的每个元素应用函数func,返回新的RDD。例如:

rdd = sc.parallelize([1, 2, 3])
result = rdd.map(lambda x: x * 2)

结果为[2, 4, 6]

filter(func) 筛选满足条件的元素。例如:

result = rdd.filter(lambda x: x > 1)

结果为[2, 3]

flatMap(func) 类似map,但func返回的是可迭代对象,结果会被扁平化处理。例如:

rdd = sc.parallelize(["hello world", "spark"])
result = rdd.flatMap(lambda x: x.split(" "))

结果为["hello", "world", "spark"]

groupByKey() 对键值对RDD按键分组。例如:

rdd = sc.parallelize([("a", 1), ("b", 2), ("a", 3)])
result = rdd.groupByKey().mapValues(list)

结果为[("a", [1, 3]), ("b", [2])]

reduceByKey(func) 对相同键的值进行聚合。例如:

result = rdd.reduceByKey(lambda a, b: a + b)

结果为[("a", 4), ("b", 2)]

行动算子(Actions)

collect() 将RDD所有元素返回给驱动程序。例如:

result = rdd.collect()

count() 返回RDD中元素的数量。例如:

result = rdd.count()

reduce(func) 通过func聚合RDD中的元素。例如:

result = rdd.reduce(lambda a, b: a + b)

saveAsTextFile(path) 将RDD保存到指定路径的文本文件中。例如:

rdd.saveAsTextFile("output_path")

高级算子

join(otherRDD) 对两个键值对RDD进行内连接。例如:

rdd1 = sc.parallelize([("a", 1), ("b", 2)])
rdd2 = sc.parallelize([("a", 3), ("c", 4)])
result = rdd1.join(rdd2)

结果为[("a", (1, 3))]

union(otherRDD) 合并两个RDD。例如:

result = rdd1.union(rdd2)

结果为[("a", 1), ("b", 2), ("a", 3), ("c", 4)]

distinct() 去重RDD中的元素。例如:

rdd = sc.parallelize([1, 2, 2, 3])
result = rdd.distinct()

结果为[1, 2, 3]

性能优化算子

persist(storageLevel) 将RDD缓存到内存或磁盘,避免重复计算。例如:

rdd.persist()

repartition(numPartitions) 调整RDD的分区数,优化并行度。例如:

rdd.repartition(4)

注意事项

  • 转换算子具有惰性求值特性,只有在行动算子触发时才会执行。
  • 窄依赖算子(如map、filter)效率高,宽依赖算子(如groupByKey)可能引发Shuffle。
  • 合理使用持久化(persist)可以提升性能,但需注意内存开销。

更多推荐