简介:这是一套面向大数据初学者的编程案例包,以Apache Spark框架为背景,系统演示如何通过Python的PySpark接口完成数据的加载、清洗、转换与聚合统计,适合正在学习Spark教程、准备大数据入门或希望快速上手分布式计算的开发者。压缩包内共十个文件,文件类型涵盖文本说明、PDF文档、VBS与Shell辅助执行脚本、Jar工具包以及License许可文件,既包含可直接运行的代码示例,也配有用于解释原理的文档和辅助脚本,整个压缩包约五百六十MB。目前已有三百六十二人参与学习。案例内容由浅入深地组织:先讲解如何创建SparkContext并读取本地或HDFS上的数据源,接着使用flatMap、filter、map等算子完成分词与过滤,再通过countByValue、reduce等操作实现词频统计和数值聚合,最后引入SparkSession与DataFrame完成结构化查询、分组计数等高级分析。跟随示例逐步运行,读者既能掌握RDD弹性分布式数据集的核心编程思想,也能学会DataFrame API的使用技巧,并理解分布式任务在本地集群中的执行逻辑,为后续应对真实大数据项目打下扎实基础。
1. 用PySpark跑通第一个大数据案例:这份代码包里有什么
拿到这份名为Python代码案例的压缩包时,我以为只是又一个脚本合集,打开目录才发现,里面全是一套面向Spark教程的PySpark示例代码。它要解决的问题很具体:让只写过单机Python的人,快速建立起“用pyspark模块操作大数据”的完整手感。适用人群主要两类——刚看完Spark理论、想让代码在自己机器上跑起来的新手,以及被业务逼着从Scala迁到PySpark的老手。这份代码最值得学习的地方,是把RDD的创建、转换、聚合和DataFrame的查询串在了一条完整链路上,而不是像官方文档那样把每个API孤立地铺开。我会按这份Spark教程的实践顺序,把环境、API和排错串起来讲,最后落到集群上可以复现的检查习惯。
2. 环境与入口选型:从SparkContext到SparkSession的演进
围绕这份代码案例,第一步并不是急着写业务逻辑,而是先理解你手里的PySpark入口从哪来、选哪个。项目中反复出现的SparkContext、SparkSession与setMaster参数,决定了你的代码是跑在本机还是集群,也决定了后续所有API的调用方式。
2.1 为什么SparkContext是旧入口,SparkSession才是当前首选
SparkContext是Spark 1.x时代的根入口,所有RDD的创建和并行计算都从它开始。在PySpark里,一行代码就能拿到连接:
from pyspark import SparkConf, SparkContext conf = SparkConf().setAppName("Spark Python Example").setMaster("local[4]") sc = SparkContext(conf=conf)这里setAppName("Spark Python Example")是给应用起名,用于在Spark UI的Application列表里区分不同任务;setMaster("local[4]")表示本地以4个线程模拟4个并行度。如果只写"local",则只有一个线程在跑,读文件、map、reduce全部串行,新手跑案例时容易误判性能,我一般至少给"local[2]"。
但如果你照旧代码学,会发现新项目已经很少直接创建SparkContext了。Spark 2.0之后引入了SparkSession,把SparkContext、SQLContext、HiveContext合并成一个统一入口。摘要描述里提到,这份代码案例既操作RDD又操作DataFrame,这就意味着你必须走SparkSession,而不是老式的“SparkContext + SQLContext”双上下文组合。继续直接创建SparkContext的不便之处在于,处理结构化数据时要额外初始化SQLContext,一旦两个上下文的配置不一致,很容易出现Hive和DataFrame数据源连不上的怪问题。
正确的打开方式是用SparkSession作为统一入口:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("Spark Python Example") \ .master("local[4]") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate()builder模式是SparkSession的推荐构造方式。appName和master的用法与SparkConf一致;config用于设置运行时参数,这里把shuffle分区数从默认的200改成了8,本地调试数据量不大时,200个空分区只会白白增加调度开销。如果你需要操作RDD,用spark.sparkContext就能拿到底层的SparkContext,两类API并存于同一个会话内。
2.2 本地与集群模式的master参数选型
代码从笔记本搬到集群,第一件要改的就是master参数。这份教程代码里的setMaster("local")只是最小可用配置,实际部署时你需要知道master各个取值对应的运行场景:
| master值 | 含义 | 适用场景 |
|---|---|---|
| local | 本地单线程,纯调试 | 第一次跑通案例 |
| local[4] | 本地4线程并行 | 单机多核验证逻辑 |
| spark://host:7077 | 连接Standalone集群 | 测试环境 |
| yarn | 提交到YARN集群 | 生产环境最常用 |
| k8s://host:443 | 提交到Kubernetes | 云原生环境 |
我的建议是,跟这份Spark教程走的时候先在local[4]下把代码跑通,再切换到yarn。因为yarn模式下,资源申请、Executor启动、日志汇聚都会变化,直接上集群排错难度陡增。本地模式适合验证RDD算子和DataFrame查询逻辑是否正确,集群模式适合验证资源参数和并行度设计是否合理。
2.3 环境变量与依赖配置:装好PySpark只是第一步
新手最容易卡在环境上,而不是代码上。PySpark 3.x要求Python 3.8以上,如果你本机还在用Python 2.7,第一步不是pip install,而是先把python安装到合适版本。确认Python版本无误后,再安装PySpark:
pip install pyspark验证安装是否可用:
python -c "import pyspark; print(pyspark.__version__)"这条命令能从Python模块层面确认PySpark是否真正进入当前环境。很多“import pyspark报错”的问题,实际上是pip装到了别的Python解释器上,比如直接用python3 -c检查过一次,再确认pip对应的解释器版本是否一致。
本地环境还要配置两个变量,我一般写在~/.bashrc里:
export SPARK_HOME=/opt/spark export PYSPARK_PYTHON=/usr/bin/python3SPARK_HOME指向Spark安装目录,PYSPARK_PYTHON告诉集群Worker使用哪个Python解释器。这个变量一旦漏配,集群环境的Executor启动时就会因为找不到Python模块而失败,而且报错不直接,经常要翻到Executor日志末尾才看得见。
3. RDD转换与聚合:把词频统计这个经典案例跑细
理解SparkContext之后,就可以正式动RDD了。摘要描述里给了一段词频统计的雏形,我这里把它扩展成完整可运行的代码串,并把每个算子的语义边界讲透。这份教程代码的核心价值就在于,它把Spark最常用的一组RDD算子按“加载—转换—聚合”串成了一条完整链路。
3.1 textFile加载:路径、分区与文件格式的取舍
RDD的数据加载从textFile开始,如果入口用的是SparkSession,代码是这样:
sc = spark.sparkContext data = sc.textFile("hdfs://path/to/your/data.txt", minPartitions=10)第一个参数是路径,支持本地文件、HDFS、S3。本地调试直接写"file:///home/user/data.txt",也可以写相对路径"data.txt";生产环境通常指到HDFS。第二个参数minPartitions是建议分区数,如果不传,Spark会按文件块大小推断。分区数决定了后续map任务的并行度——分区太少,资源闲置;分区太多,调度开销反而把时间吃掉。
注意,textFile是按行读取的,每一行成为RDD中的一个元素。如果源文件是CSV带表头,读进来之后需要额外处理第一行;如果源文件是JSONL(每行一个JSON对象),则可以用map配合json.loads逐行解析。这份代码案例里用的是文本文件,所以没有涉及这些格式转换,但你如果要套用到业务数据,这个边界得先想清楚。
3.2 map、flatMap与filter:这三个算子的语义差别
词频统计的标准拆解是这一段:
words = data.flatMap(lambda line: line.split()) filtered_words = words.filter(lambda word: word.startswith("a"))flatMap与map的区别是新人最容易翻车的地方。map对一个元素只能产出一个元素,flatMap允许一个元素产出多个元素。这里一行文本被line.split()拆成了若干个单词,所以必须用flatMap。如果你误用map,RDD里的每条记录会变成一个数组,后续的filter和countByValue都会对着数组操作,结果完全不对。
filter则是对每个元素做布尔判断,保留结果为True的元素。它在RDD转换中不会改变元素个数之外的结构,只是把不需要的记录丢弃。这三个算子组合起来,正好覆盖了“拆分—过滤—统计”的常见数据处理节奏。
补上统计和输出,代码变成完整可跑的版本:
words = data.flatMap(lambda line: line.split()) a_words = words.filter(lambda word: word.startswith("a")) for word, count in a_words.countByValue().items(): print(word, count)这里countByValue()是一个action操作,它会真正触发计算,而不是像flatMap和filter那样构建DAG后惰性等待。Spark的惰性机制是新手最容易误判的一点——写了flatMap就以为已经开始跑了,其实它只是往DAG里加了一个节点,直到遇到action才真正开始分配任务执行。
3.3 聚合操作:reduce、countByValue与groupBy的边界
词频统计还有另一种更常见的实现,用reduceByKey:
pairs = words.map(lambda word: (word, 1)) word_counts = pairs.reduceByKey(lambda a, b: a + b) for word, count in word_counts.collect(): print(word, count)reduceByKey接收一个二元函数,把相同key的value两两合并。与groupBy相比,reduceByKey会把聚合逻辑前移到map端做combine,shuffle量小得多,是词频统计的首选方案。groupBy则是把所有value原封不动汇聚到一个迭代器里,适合分组之后还需要完整列表的场景,但shuffle数据量明显更大。
countByValue则更特殊——它直接按元素值统计,返回的是Python原生字典,不是RDD。对于去重统计很方便,但如果数据量很大,全部结果堆在Driver内存里会OOM。它的适用边界就在这:小数据量、需要快速看结果时用;大数据量还是走reduceByKey,把结果保留在分布式数据结构里。
如果要做TopN,可以继续串一个sortBy:
sorted_counts = word_counts.sortBy(lambda x: x[1], ascending=False).take(10) for word, count in sorted_counts: print(word, count)sortBy的第一个参数是排序键函数,这里按count值排序;take(10)只取前10条,比collect()把所有结果拉到本地更安全。我一般会在调试时先用take(20)看数据结构,确认没问题再全量collect()。
4. DataFrame与SQL:结构化查询的迁移与参数坑
RDD能解决大部分非结构化数据处理问题,但一旦源数据有明确的行列结构,更高效的做法是转成DataFrame。摘要描述里的CSV读取示例,正好引出了字段类型推断和SQL查询这两个关键点。
4.1 从RDD到DataFrame的两条路线
如果你已经有一套RDD,想转成DataFrame去执行SQL级查询,有两条路:
df_from_rdd = rdd.toDF(["word", "count"]) from pyspark.sql.types import StructType, StructField, StringType, LongType schema = StructType([ StructField("word", StringType(), True), StructField("count", LongType(), True) ]) df_from_rdd2 = spark.createDataFrame(rdd, schema)toDF(["word", "count"])适合快速转换,列名通过列表传入,字段类型按值自动推断;createDataFrame(rdd, schema)则允许完全控制字段名和类型。两者在Spark内部的执行路径几乎一致,区别在于类型控制粒度。业务数据里经常遇到“某个字段被推断成了string,后续聚合报类型错误”,所以只要能事先确定schema,我一般直接用第二种写法,省得后面再修。
DataFrame的API风格和RDD差别很大,称得上两套思维。RDD偏函数式,DataFrame偏关系型。筛选以“a”开头的单词,DataFrame这样写:
df_filtered = df.filter(df.word.startswith("a")) df_count = df_filtered.groupBy("word").count() df_count.show()filter里用的是Column表达式,df.word.startswith("a")会生成一个布尔Column,groupBy("word").count()返回一个新的DataFrame,show()在控制台打印前20行。这套API对传统SQL开发者的迁移成本很低,也是这份Spark教程特意把DataFrame和RDD放在一起讲的原因。
4.2 CSV读取参数:inferSchema与header的实际行为
摘要描述里有一段CSV读取:
df = spark.read.csv("hdfs://path/to/csv", inferSchema=True, header=True) result = df.groupBy("column_name").count()inferSchema=True表示自动推断每列类型,header=True表示把第一行当列名。这两个参数单独看都不难,但组合起来有隐性成本。下表是spark.read.csv最常用参数:
| 参数 | 默认值 | 作用 |
|---|---|---|
| header | false | 是否把第一行当列名 |
| inferSchema | false | 是否自动推断列类型 |
| sep | , | 列分隔符 |
| mode | PERMISSIVE | 损坏记录容错模式 |
| samplingRatio | 1.0 | 类型推断时的采样比例 |
inferSchema默认采样全部数据,文件很大时这一步会拖慢整体读取;如果把samplingRatio调小,推断结果又会失真,常见坑是把整数列推断成字符串。更稳的做法是手动定义schema,跳过推断阶段:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType custom_schema = StructType([ StructField("city", StringType(), True), StructField("sales", IntegerType(), True) ]) df = spark.read \ .option("header", "true") \ .option("sep", ",") \ .schema(custom_schema) \ .csv("hdfs://path/to/csv") result = df.groupBy("city").sum("sales") result.show()手动指定schema之后,读取阶段不再做类型推断,速度更快,类型也不会漂移。groupBy("city").sum("sales")按城市分组后求销售额总和,结果仍然是DataFrame,show()打印结果。
4.3 用SQL查询替代链式调用:什么时候值得
DataFrame还有一种更接近关系型数据库的操作方式,注册成临时视图后直接写SQL:
df.createOrReplaceTempView("sales_tbl") sql_result = spark.sql(""" SELECT city, SUM(sales) AS total_sales FROM sales_tbl GROUP BY city """) sql_result.show()createOrReplaceTempView注册的是当前SparkSession内的临时视图,会话结束即失效。对于多表join、嵌套子查询这类逻辑,SQL的可读性明显高于链式API,团队里如果SQL背景的人多,用Spark SQL迁移起来最平滑。要注意的是字段名里如果含空格或SQL保留字,需要在SQL里用反引号包裹,这个细节经常把人绊一下。
什么时候不值得用SQL?当查询逻辑已经被前面的DataFrame步骤加工过,且只做简单聚合时,链式API更简洁,还能直接复用已经构建好的Column表达式。两种风格各有适用场景,这份教程代码把两种都覆盖了,练习时可以分别跑一遍,对比可读性和执行计划。
5. PySpark避坑指南:本地跑通到集群提交的五个坎
在本地IDE里跑通案例不等于在生产环境能跑通。从local[4]到yarn,这中间隔着Python环境、内存、序列化、路径四类问题。下面这五个坑,是这份Spark教程落地过程中一定会遇到的坎,每个按“现象→原因→解决”的记录方式拆开。
5.1 Executor崩溃:Python版本不匹配
现象:本地IDE跑得好好的,spark-submit提交到集群后,Executor启动即退出,日志里有Python in worker has different version或ModuleNotFoundError: pyspark。
原因:Driver端的Python环境是本地解释器,Executor端默认使用系统的python,两边版本不一致或site-packages路径不同,导致Executor进程里import不到pyspark模块。
解决:在提交命令里显式固定解释器路径:
spark-submit \ --master yarn \ --conf spark.pyspark.python=/usr/bin/python3 \ --conf spark.pyspark.driver.python=/usr/bin/python3 \ app.pyspark.pyspark.python控制Worker端解释器,spark.pyspark.driver.python控制Driver端解释器。集群环境里最好统一用绝对路径,不要依赖PATH里的默认值。我习惯在提交脚本里把这两个配置固定写死,防止不同节点默认python版本漂移。
5.2 Shuffle阶段OOM:数据倾斜还是分区太少
现象:执行reduceByKey或groupBy时,某个Executor被标记为Container killed,日志尾部GC时间异常高。
原因:少数key的数据量远大于其他key,单个分区的reduce任务处理压力过大;或者spark.sql.shuffle.partitions仍为默认200,但单分区数据量远超内存承载。
解决:先调整分区数:
spark.conf.set("spark.sql.shuffle.partitions", "500")分区数调大之后如果还挂,就要怀疑数据倾斜。常见做法是对大key加盐,把大key拆成多个子key,聚合完再去盐合并。加盐的具体逻辑依赖业务,但事前用countByKey看一眼key分布,是值得养成的习惯。
5.3 lambda序列化错误:算子内部别引外部对象
现象:运行到某个map或foreach算子时报PicklingError或Could not serialize object。
原因:lambda函数里捕获了不可序列化的外部对象,比如数据库连接、打开的文件句柄、SparkSession本身。PySpark会把算子函数分发到各Executor,捕获的大对象必须能被pickle,否则直接报错。
解决:把对象创建移到算子内部,或改成模块级具名函数:
def process_row(row): # 连接在函数内部创建,每个分区创建一次 return transform(row) rdd.map(process_row)模块级函数天然比lambda更可预测,也能避免隐式捕获。曾经遇到一个案例,开发者在flatMap里直接调用了外层函数里的spark.sql,运行到一半就序列化失败,排查了很久才找到是spark这个上下文被闭包捕获了。
5.4 inferSchema把数值列读成字符串
现象:CSV读取后,某列在df.printSchema()里显示为string,但肉眼可见全是整数。
原因:inferSchema只采样部分行做推断,如果采样区间内数据恰好都是空值或带引号字符串,类型推断就会出错。samplingRatio被调小时,这种错误更容易出现。
解决:不要在生产链路里依赖inferSchema,手动声明StructType。第四章已经给出完整代码。如果只是本地调试,可以临时把samplingRatio调大重跑,但不要把这个习惯带到集群任务里。
5.5 本地能跑集群报FileNotFoundError
现象:spark-submit提交后,作业启动即报FileNotFoundError: [Errno 2] No such file or directory。
原因:代码里用了本机绝对路径读取资源文件,集群Executor节点上没有这个路径。
解决:数据文件放HDFS或对象存储,小配置文件通过SparkFiles.add分发:
from pyspark import SparkFiles spark.sparkContext.addFile("config.json") def load_config(): import json with open(SparkFiles.get("config.json"), "r") as f: return json.load(f)SparkFiles.get会把文件拉取到Executor的工作目录,不依赖本机路径。这个坑几乎所有上过集群的人都踩过,属于PySpark从本地迁移到集群的血泪经验。
6. 用Spark UI验证调优:缓存与广播变量的检查流
代码能在集群稳定运行之后,下一步是验证调优效果。Spark UI里能看到的信息很多,但核心只需要关注两个地方:Job页面的DAG图和Storage标签页。
### 6.1 通过DAG看Stage划分
Spark作业在UI里会显示成DAG图。每次shuffle都会切分出一个新的Stage,Stage数量越多,shuffle开销越大。点进每个Stage能看到Shuffle Read和Shuffle Write的字节数,如果Write量特别大,说明上游算子产生了过多中间数据,优先考虑把map端的聚合需求合并进reduceByKey这类算子,减少落盘量。
### 6.2 缓存与广播变量的选择
同一个RDD如果被多个action复用,每次都会从头计算。加缓存能直接消除重复计算的开销:
base_rdd = data.flatMap(lambda line: line.split()).cache() first = base_rdd.count() second = base_rdd.filter(lambda word: len(word) > 3).count()不加cache()时,第二次count仍会触发一次完整的DAG计算;加缓存后,第二次直接在内存里复用。缓存级别默认为MEMORY_ONLY,数据量大的场景改用MEMORY_AND_DISK,防止内存不够时直接落盘重算。
广播变量则适合把只读配置下发到每个Executor,避免在算子里反复访问外部系统:
lookup_dict = spark.sparkContext.broadcast({...}) rdd.map(lambda row: (row, lookup_dict.value.get(row)))broadcast.value只读不能改,适合承载字典、配置表这类全局数据,尤其当它被多个map算子引用时,能明显减少传输开销。
### 6.3 验证方法
调优有没有效果,不能靠感觉。跑完作业后切到Spark UI的Storage标签页,确认缓存数据确实存在于内存;再对比两次作业的耗时差,就可以判断缓存是否值得。从那以后我每次提交PySpark作业,都会强制走一遍这个检查流:先看DAG的Stage数量和shuffle量,再决定要不要加cache或广播变量,最后跑完核对Storage标签页和总耗时。希望帮到你。
本文还有配套的精品资源,点击获取