尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

PySpark实战:从单机脚本到分布式数据处理的完整指南

PySpark实战:从单机脚本到分布式数据处理的完整指南 简介围绕Apache Spark的Python编程示例这份rar压缩包面向希望借助PySpark进入大数据处理领域的初学者与开发者。压缩包共10个文件txt说明文档用于代码解读和环境配置提示pdf阅读材料提供基础理论梳理vbs与sh脚本帮助在Windows/Linux下快速准备Spark运行环境jar组件与license授权文件则保证工具链的完整性整体大小559.17MB适合离线学习与对照运行。案例从SparkContext创建与textFile加载数据出发逐步演示flatMap拆分单词、filter过滤指定前缀、countByValue词频统计等RDD转换与聚合操作同时也基于SparkSession演示DataFrame读取CSV、按列分组计数的结构化数据处理方法。这些代码覆盖了Spark常用API的典型调用路径和关键配置方法能帮助读者理解RDD与DataFrame两种编程模型并积累可直接复用的PySpark脚本。已有362人学习下载对正在准备大数据开发入门或需要快速上手Spark的读者来说是一套兼具教学与参考价值的配套案例。1. PySpark代码案例实战从单机脚本到分布式数据处理的一步到位有一次我拿公司线上商城的一批全量访问日志做分析单机 Python 脚本跑一遍清洗加聚合要小半天数据量一上来内存直接被打满。后来把同样逻辑改写成 PySpark 代码加载、过滤、分组统计三步走十分钟内出结果这套 Python 代码案例包演示的就是这条路。它不教 Python 语法而是把你已经会的 Python 迁移到 Spark 分布式计算模型上覆盖 SparkContext 初始化、RDD 转换、textFile 读数据、DataFrame 分组聚合这些核心环节还带着异常处理和参数调优的思路。适合已经有 Python 基础、想进入大数据开发或者正在做数据处理量升级的工程师。下面按项目推进顺序把可复现的代码和关键参数拆开讲清楚。2. SparkContext与RDD核心抽象PySpark任务的入口在这里定生死2.1 RDD的弹性来自分区与血缘RDDResilient Distributed Dataset是 Spark 里最核心的数据抽象它把数据集合切分成多个 partition 分布在集群的 executor 上每个 partition 是数据集的一个分片Spark 以 partition 为单位调度并行计算。很多人把 RDD 理解为 Python 里的 list其实两者的差别很大RDD 是不可变的且逻辑上跨机器分布普通 list 无法在集群上并行处理。RDD 的“弹性”主要体现在血缘lineage机制上每个 RDD 都记录了自己是由哪个父 RDD、经过哪种转换算子得到的一旦某个分区的数据因为节点故障丢失Spark 不需要重新读取源文件而是根据血缘关系重新执行转换计算出丢失的分区。理解了这一点你就会明白为什么 Spark 官方建议 RDD 上多做转换、少做行动操作因为转换只是在构建血缘链不真正计算结果。2.2 SparkContext参数appName、master和executor内存的关系PySpark 程序的入口是 SparkContext它负责和集群管理器通信、申请资源、创建 RDD。初始化时最常改的就是 setAppName 和 setMaster 这两个参数它们直接决定任务在集群里以什么身份运行、占多少资源。from pyspark import SparkConf, SparkContext conf SparkConf() \ .setAppName(access-log-analysis) \ .setMaster(local[*]) \ .set(spark.executor.memory, 4g) sc SparkContext(confconf) print(sc.version)这段代码里setAppName 填的是任务在 Spark Web UI 上显示的名字排查问题时要靠它定位作业建议用有业务含义的名称setMaster 指定资源调度地址本地调试填local[*]意思是使用本机所有 CPU 核并行处理local[2]则是只开 2 个线程适合快速验证逻辑。部署到集群时这里要改成yarn或者spark://ip:7077。.set(spark.executor.memory, 4g)给每个 executor 分配 4GB 内存注意这是堆内内存实际可用还要算上 overheadsubmit 脚本里超过物理内存很容易被 YARN 直接 kill 掉。2.3 textFile加载与flatMap、filter转换的配合方式创建好 SparkContext 之后第一步通常是加载数据。Spark 支持文本文件、HDFS、Cassandra、JDBC 等多种数据源最常见的仍然是用 textFile 读日志和明文数据。raw_rdd sc.textFile( hdfs://namenode:8020/logs/access_2024.log, minPartitions8 ) def parse_line(line: str): fields line.split(|) return (fields[0], fields[3]) # (user_id, action) parsed_rdd raw_rdd.map(parse_line) valid_rdd parsed_rdd.filter(lambda x: x[1] in (view, click)) print(valid_rdd.take(5))textFile 的第二个参数 minPartitions 表示期望的最小分区数数据量大的时候可以调大它来增加并行度。这里用 map 做一行数据的结构化解析map 是一对一转换每条输入产生一条输出filter 则根据布尔表达式决定记录是否保留。再看另一个场景比如对日志内容做分词words raw_rdd.flatMap(lambda line: line.split( )) filtered_words words.filter(lambda word: word.startswith(a)) word_count filtered_words.count()flatMap 与 map 的区别在于map 返回值必须是单个对象flatMap 可以把每个输入映射成多个输出然后展平非常适合单词拆分这类操作。count() 是一个行动操作它会触发实际计算并返回结果也是验证整个转换链是否正确的第一步。3. SparkSession与DataFrame结构化处理从写SQL开始3.1 SparkSession为什么成为统一入口RDD 虽然灵活但 API 粒度较细处理结构化数据时要写很多样板代码。从 Spark 2.0 开始官方把 SparkContext、SQLContext、HiveContext 整合进 SparkSession一个对象同时处理元数据、SQL 查询和 DataFrame 的创建。在 PySpark 程序中先用 SparkSession.builder 构建入口再去创建 DataFrame这也是当前数据开发的主流写法。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(user-behavior-analysis) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate()SparkSession 和 SparkContext 不需要同时显式创建SparkContext 会被 SparkSession 内部自动管理。.config(spark.sql.shuffle.partitions, 4)是一个值得注意的参数它控制 SQL 和 DataFrame 操作中 shuffle 之后的分区数默认是 200本地测试时数据量小200 个空分区只会浪费调度资源改成 4 个能让本地跑得更快。提交到集群再把该参数调回 200 或按数据量估算。3.2 CSV 加载时 header 与 inferSchema 的典型配置DataFrame 的一个优势是能借助 Catalyst 优化器自动优化执行计划但前提是它知道每一列的数据类型所以读文件时类型推断很关键。df spark.read.csv( hdfs://namenode:8020/data/user_behavior.csv, headerTrue, inferSchemaTrue ) df.printSchema()参数作用建议header是否把首行作为列名数据文件有表头时必须为 TrueinferSchema自动推断各列类型小文件可以开超大文件会额外扫描一次建议关闭并手动指定 schemaescape指定转义字符字段中包含分隔符时配置multiLine是否允许字段跨行CSV 字段带换行时设为 TrueinferSchema 在数据量大时会有额外的 IO 开销生产环境读取 TB 级 CSV 更稳的做法是先用小样本推断出类型然后手动定义 StructType既准确又省去一次额外扫描。3.3 列选择、过滤与临时视图下的 SQL 查询DataFrame 的 API 写法比 RDD 更接近 SQL 思维。常见操作是先做列投影和过滤再交给 Spark 执行统计。clicks df.select(user_id, action, city_id) \ .filter(df[action] click) \ .filter(df[city_id].isNotNull()) clicks.show(10)select 负责选出需要的列作用是尽早裁剪数据列减少后续计算中的数据量。filter 对应 SQL 中的 whereSpark 会把多个 filter 合并并下推到读取阶段。另一个实用操作是把 DataFrame 注册成临时视图直接用 SQL 做复杂查询df.createOrReplaceTempView(behavior) result_sql spark.sql( SELECT action, count(*) AS cnt FROM behavior GROUP BY action ) result_sql.show()createOrReplaceTempView 创建的是会话级临时表SparkSession 停止后自动销毁不影响其他任务。用 SQL 写分组统计比链式 API 更容易阅读复杂 join 场景下也更好维护。这里可以看到 DataFrame 和 RDD 的根本区别DataFrame 带 schema 信息Spark 能对执行计划做优化比如列裁剪和谓词下推而 RDD 不知道列类型只能按 Java/Python 对象做序列化处理。4. 聚合计数的性能取舍countByValue、reduceByKey与groupByKey怎么选4.1 聚合API的数据流向差异做单词计数或分类统计时PySpark 有好几个聚合入口但它们的行为差异很大用错了会直接影响任务是否撑得住。我用一张表来说明它们的主要区别。聚合方式返回类型数据流向适用场景countByValuePython dict全量结果回收到 driver结果集很小快速验证reduceByKeyRDDmap 端预聚合再 shuffle 局部结果大规模 key 求和、计数groupByKeyRDD每个 key 的完整 value 列表跨节点传输需要对全量 value 做非加法计算groupBy countDataFrameCatalyst 优化执行SQL 风格统计分析很多新手以为 groupByKey 是标准聚合方式其实它把所有值原样送到下游节点完全不做本地合并例如对 1 亿条数据按 key 分组key 只有 100 个groupByKey 依然会传输 1 亿条记录而 reduceByKey 先在每个分区内做一次 add再只把局部汇总结果 shuffle 出去网络开销小一个量级。4.2 reduceByKey的map端combine及与groupByKey的代码差异用一段代码直观展示两者的差别场景是对上一步解析出的 action 字段做计数。from operator import add pairs valid_rdd.map(lambda x: (x[1], 1)) counts_rdd pairs.reduceByKey(add) counts_rdd.saveAsTextFile(output/action_counts) # 不推荐写法仅用于对比 wrong_way pairs.groupByKey() \ .mapValues(lambda vs: sum(vs))map 把原始数据变成(action, 1)键值对reduceByKey 先对每个分区内部相同 key 调用 add 进行局部相加再把局部结果按 key 合并整个过程中每个 key 的中间结果都是一个小整数。groupByKey 则要先把每个 key 对应的所有 value 组装成迭代器再传输给下游节点这里传输的是完整列表占用网络带宽和序列化时间。数据量越大这两种写法的差距越明显。如果非要在 RDD 上做聚合建议优先考虑 combineByKey 或者 reduceByKey。另外注意 countByValue 在 RDD 上直接调用时会把整个结果集拉回 driver分区数多、key 数量大时极容易内存溢出它更适合在采样后的小数据集上做分布探查不要直接用在全量生产数据上。4.3 广播变量解决小表关联大表的性能问题实际业务里经常要把 RDD 或 DataFrame 里的 id 字段关联成名称字段比如把 city_id 转成城市名。常见做法是写一个 join但如果关联表很小join 反而会引入一次全量 shuffle。更优的方案是使用广播变量。city_map {101: 北京, 102: 上海, 103: 广州} bc_city sc.broadcast(city_map) def add_city_name(row): user_id, action, city_id row return (user_id, action, bc_city.value.get(city_id, unknown)) df_rdd clicks.rdd.map(add_city_name) df_rdd.take(5)sc.broadcast 把 city_map 序列化后分发到每个 executor 的内存里executor 上的任务直接从本地读取这份只读字典不需要跨节点传输。使用时注意两点第一广播变量只读不要在 executor 端尝试修改第二只适合小数据量几 MB 到几十 MB 的维度表广播效果最好几百 MB 的字典反而会给每个 executor 增加内存压力这时还是应该用 join。5. 缓存分区与任务异常定位PySpark任务跑挂后先查这几处5.1 cache与persist的触发时机同一个 DataFrame 如果会被多个行动操作反复使用就应该缓存起来否则每个 action 都会从头计算整条血缘链。df.cache() df.count() result df.groupBy(action).count() result.show()cache 的默认存储级别是 MEMORY_AND_DISK内存不够时溢出到磁盘不会直接报错。也可以用 persist 手动指定级别数据量大且需要跨多个 action 复用时选择 MEMORY_AND_DISK_SER 能减少内存占用代价是增加反序列化时间。缓存是懒执行的第一次 count() 才会真正把数据存入存储内存之后 groupBy 查询时就直接读缓存了。任务结束后调用 unpersist() 释放。5.2 分区数调整与并行度任务跑得慢或者并发度偏低通常要先看一眼分区数。通过 Spark UI 的 Stage 页能看到当前 RDD 的分区数和每个分区处理的数据量再决定用 coalesce 还是 repartition。merged parsed_rdd.coalesce(4) # 减少分区不触发 shuffle expanded merged.repartition(16) # 增加分区触发 shufflecoalesce 只做分区合并不会产生网络传输适合在 shuffle 之后合并小文件repartition 会重新划分数据伴随一次 shuffle适合在数据倾斜时增加分区缓解单节点压力。SQL 层面的分区数则由 spark.sql.shuffle.partitions 控制我通常先看到 Spark UI 里单个 stage 的输入数据量再按每个分区 64MB 到 128MB 来估算合理分区数。5.3 从Spark UI定位OOM和Executor被Kill任务执行失败时第一件事不是调代码而是打开 Spark Web UI此时先看 Executors 页面关注 Storage Memory 使用量和 Shuffle Spill 数值。Shuffle Spill 表示数据在 shuffle 过程中溢写到磁盘如果这个值特别大说明 executor 内存严重不足。再看 Event Timeline 里失败任务的报错信息如果是Container killed by YARN for exceeding memory limits优先调整 executor 内存配置比如在提交脚本中加--executor-memory 8g --executor-cores 4并同步调整spark.memory.fraction这个参数决定统一内存区域中执行和存储的比例默认 0.6 不适合大量 shuffle 计算。如果代码里做了 collect()还要检查是否把全量数据拉回了 driver这个操作比调参更容易把任务拖垮。本文还有配套的精品资源点击获取
返回列表