Spark是大数据领域主流分布式计算框架,核心抽象为弹性分布式数据集(RDD),分区机制与算子操作是其高效计算的关键。本文精简梳理RDD特性、分区规则、高阶函数及常用算子,搭配实战代码,帮助初学者快速掌握核心用法、避开常见坑点。

一、RDD的诞生:为什么需要弹性分布式数据集?

传统集合(如Java ArrayList、Python List)仅能在单台服务器内存中存储,无法实现分布式存储与计算,难以满足大数据处理需求。

示例代码(读取外部文件并处理):

# step1: 读取数据
input = sc.textFile("输入路径")
# step2: 处理数据(map、filter等操作)
result = input.map(lambda x: x.strip())
# step3: 保存结果
result.saveAsTextFile("输出路径")

Spark设计的弹性分布式数据集(RDD),可像集合一样存储元素,同时具备分布式特性,支持跨节点并行计算,解决了传统集合的局限。

1.1 RDD的定义与创建方式

官方定义:RDD是容错、可并行操作的分布式元素集合,核心创建方式有两种:

  • 并行化集合:通过parallelize方法将Driver端集合转为RDD,指定分区数实现分布式存储。

  • 读取外部存储:通过textFile、wholeTextFile等方法,读取HDFS、本地文件等外部数据转为RDD。

实战代码示例:

import os
from pyspark import SparkContext, SparkConf

if __name__ == '__main__':
    # 配置环境变量(按自身环境修改路径)
    os.environ['JAVA_HOME'] = 'C:\Program Files\Java\jdk1.8.0_241'
    os.environ['HADOOP_HOME'] = 'D:\hadoop-3.3.1'
    os.environ['PYSPARK_PYTHON'] = 'C:/ProgramData/Miniconda3/python.exe'
    
    # 初始化SparkContext
    conf = SparkConf().setMaster("local[*]").setAppName("RDD_Create_Demo")
    sc = SparkContext(conf=conf)
    
    # 方式1:并行化集合创建RDD
    data = [1,2,3,4,5,6,7,8,9,10]
    list_rdd = sc.parallelize(data, numSlices=2)  # numSlices指定分区数
    list_rdd.foreach(lambda x: print(x))
    
    # 方式2:读取外部文件创建RDD
    file_rdd = sc.textFile("../datas/filter.txt", minPartitions=2)  # 最小分区数
    file_rdd.foreach(lambda line: print(line))
    
    sc.stop()  # 关闭SparkContext

1.2 RDD的五大核心特性

RDD的五大特性是其容错性和并行性的基础,也是理解Spark计算模型的关键:

  1. 分区特性:由多个分区组成,每个分区存储在不同节点,是分布式计算的最小单元。

  2. 并行转换特性:转换操作(如map)对所有分区并行执行,每个分区对应一个Task。

  3. 依赖关系(血脉机制):保存与其他RDD的依赖关系,数据丢失可通过依赖重新构建。

  4. 可选分区器:KV类型RDD可自定义分区器,默认提供HashPartition和RangePartition。

  5. 本地优先计算:优先将Task分配到数据所在节点,减少数据传输,提升效率。

二、RDD分区的设定规则:影响计算效率的关键

分区数直接影响并行度和计算效率,不同创建方式的分区规则如下:

2.1 parallelize创建RDD的分区规则

  • 指定分区数:numSlices参数直接决定分区数。

  • 未指定分区数:由spark.default.parallelism参数决定,不同运行模式默认值不同。

2.2 textFile读取外部数据的分区规则

  • 未指定分区数:取spark.default.parallelism和2的最小值。

  • 指定分区数:为最小分区数,实际分区数按HDFS分片规则(默认128M/片)计算。

2.3 不同运行模式的默认并行度

  • Local模式:等于本地CPU核数(local[*]自动适配)。

  • Mesos模式:默认并行度为8。

  • 集群模式(Standalone、Yarn):取所有Executor总核数与2的最大值。

注:子RDD默认分区数与父RDD一致,可通过repartition、coalesce算子手动修改。

三、高阶函数及Lambda表达式:Spark算子的基础

Spark算子依赖高阶函数和Lambda表达式,Python中Lambda可简化代码、提升开发效率。

3.1 Lambda表达式基础

Lambda是匿名函数,适用于函数仅使用一次、逻辑仅一行的场景,语法:lambda 参数: 表达式。

import math

list1 = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
# Lambda实现平方、立方计算
pingFang = lambda x: x * x
liFang = lambda x: math.pow(x, 3)

print(list(map(pingFang, list1)))  # 平方结果
print(list(map(liFang, list1)))   # 立方结果

3.2 高阶函数定义

高阶函数是参数包含另一个函数的函数,Spark的map、filter等算子均为高阶函数。

关键区别:函数本地计算,无法并行;算子分布式计算,可并行处理RDD分区数据。

四、RDD常用基础算子:实战必备

Spark算子分为两类:转换算子(Lazy加载,返回RDD)和触发算子(触发Job,返回非RDD)。

4.1 算子分类速查表

转换算子:map、filter、flatMap、reduceByKey、groupByKey、sortBy、sortByKey、union、distinct等。

触发算子:count、foreach、saveAsTextFile、first、take、collect、reduce、top等。

Shuffle算子(影响性能):reduceByKey、groupByKey、sortBy、join、distinct等。

4.2 常用转换算子实战

(1)map算子:一对一转换

对RDD每个元素执行函数,返回新RDD,适用于一对一转换。

# 计算每个元素的立方
list01 = [1,2,3,4,5,6]
listRdd = sc.parallelize(list01)
mapRdd = listRdd.map(lambda x: math.pow(x, 3))
mapRdd.foreach(lambda x: print(x))

(2)flatMap算子:扁平化转换

展开嵌套集合,类似SQL的explode函数。

# 按“/”分割每行歌曲,展开为单个歌曲名
fileRdd = sc.textFile("../datas/a.txt", 2)
flatRdd = fileRdd.flatMap(lambda line: line.split("/"))
flatRdd.foreach(lambda x: print(x))

(3)filter算子:数据过滤

按条件过滤元素,适用于行级过滤(类似SQL的where)。

import re

# 过滤状态不为-1且字段数为4的行
fileRdd = sc.textFile("../datas/b.txt", 2)
filterRdd = fileRdd.filter(lambda line: re.split(r"\s", line)[2] != '-1' and len(re.split(r"\s", line)) == 4)
filterRdd.foreach(lambda x: print(x))

(4)sortBy与sortByKey:排序算子

均实现全局排序,sortBy适用于所有RDD,sortByKey仅适用于KV类型RDD。

注:直接用foreach打印可能顺序错乱,需用collect()收集或设分区数为1。

# sortBy:按年龄降序排序
fileRdd01 = sc.textFile("../datas/c.txt")
sortByRdd = fileRdd01.sortBy(keyfunc=lambda x: int(x.split(",")[1]), ascending=False)
print(sortByRdd.collect())

# sortByKey:KV类型按Key降序排序
rdd2 = sc.parallelize([("word", 10), ("hello", 20), ("laoyan", 1)], numSlices=2)
rdd2.sortByKey(ascending=False).foreach(lambda x: print(x))

(5)groupByKey与reduceByKey:分组聚合算子

仅适用于KV类型RDD,reduceByKey先分区内聚合再Shuffle,性能优于groupByKey。

# 统计单词出现次数
rdd5 = sc.parallelize([("word", 10), ("word", 5), ("hello", 100), ("hello", 20)], numSlices=3)

# groupByKey实现(不推荐)
rdd5.groupByKey().map(lambda x: (x[0], sum(x[1]))).foreach(lambda x: print(x))

# reduceByKey实现(推荐)
rdd5.reduceByKey(lambda x, y: x + y).foreach(lambda x: print(x))

(6)分区算子:mapPartitions与foreachPartition

对整个分区操作,减少外部连接创建次数,提升性能。

# mapPartitions:计算每个分区总和
rdd1 = sc.parallelize([1,2,3,4,5,6,7,8], 2)
rdd1.mapPartitions(lambda x: [sum(x)]).foreach(print)

# foreachPartition:遍历每个分区
def showMsg(x):
    for i in x:
        print(i)
rdd1.foreachPartition(showMsg)

4.3 常用触发算子实战

  • count:统计元素个数;foreach:遍历无返回值;saveAsTextFile:保存数据到外部系统。

  • collect:汇总数据到Driver内存(大数据量慎用);reduce:聚合计算。

# 触发算子综合示例
rdd1 = sc.parallelize([1,2,3,4,5,6,7,8,9,10], 2)
print(rdd1.count())  # 统计个数
rdd1.foreach(lambda x: print(x*10))  # 遍历打印
rdd1.saveAsTextFile("../../output111")  # 保存结果
print(rdd1.collect())  # 收集打印
print(rdd1.reduce(lambda tmp, item: tmp + item))  # 求和

五、常见问题与注意事项

  1. sortBy/sortByKey打印错乱:用collect()收集或设分区数为1。

  2. collect内存溢出:大数据量禁用,可用take(N)获取前N条数据。

  3. Shuffle优化:减少Shuffle操作,优先用reduceByKey,合理设置分区数。

  4. 分区算子:优点是提升性能,缺点是单分区数据过大易内存溢出。

六、总结

本文精简梳理了Spark核心知识点,涵盖RDD特性、分区规则、高阶函数及常用算子,搭配实战代码,覆盖初学者入门核心需求。Spark的核心是“分布式”和“懒加载”,合理运用分区和算子可大幅提升计算效率。

建议多动手运行代码,熟悉算子用法和坑点,后续可深入学习Shuffle机制、数据倾斜等进阶内容。

觉得有帮助可点赞收藏,后续将持续更新Spark进阶知识点!

更多推荐