PySpark原理介绍
小文件处理
背景:hive 分区如果产生了大量小文件,不仅会消耗存储元数据quota,还会导致在读取该分区时性能和效率低下,大量的时间浪费在了元数据获取上,同时在数据存储上效率也偏低,存储浪费在了元数据record上,因此小文件是很值得进行优化的功能点。
在没有shuffle阶段的处理过程中出现小文件,通过distribute by cast(rand() * 10 as int) 增加一个shuffle阶段将小文件重分区成10个分区,减少小文件。
一.SQL写入
1.单分区写入:
作用:全局随机打散,均匀分区。
适用场景:解决数据倾斜,但会产生大量小文件。
单日数据 ≤ 2GB → distribute by cast(rand() as int) – 强制将所有数据给到 key=0,无论数据大小永远1个文件。
单日数据 2GB ~ 10GB → distribute by cast(rand() * 10 as int) – 固定10份,控制份数
-- spark和hive都支持INSERTINTOTABLEtableaPARTITION(dt)SELECTcol1,col2,dtFROMtableb DISTRIBUTEBYrand();-- DISTRIBUTE BY rand() 只定义路由规则,不规定分区总数;文件数量由【引擎自动分区策略 + task参数】决定。-- spark支持,hive不支持。INSERTINTOTABLEtableaPARTITION(dt)SELECTcol1,col2,dtFROMtablebCOALESCE(1);2.Hint重分区方式(spark支持):
强制 Shuffle 重分区再收拢到 n个 分区
INSERTINTOTABLEtableaPARTITION(dt)SELECT/*+ REPARTITION(1) */col1,col2,dtFROMtableb;3.动态分区(多分区):
作用:按 dt 分区 + 分区内随机,全部进入单个 task,仅产生一个文件。
适合场景:动态分区、多日期、生产标准
单日数据 > 10GB → DISTRIBUTE BY dt, rand():分散到多个 task,负载均衡,打散热点,单个分区多文件。
-- 先开合并(防止小文件)SEThive.merge.mapfiles=true;SEThive.merge.mapredfiles=true;SEThive.merge.size.per.task=268435456;-- 256MBSEThive.merge.smallfiles.avgsize=134217728;-- 128MBINSERTINTOTABLEtable1PARTITION(dt)SELECTcol1,col2,dtFROMtable2 DISTRIBUTEBYdt,rand()-- 按分区+随机打散SORTBYdt,rand();-- 每个分区内有序,避免碎片4.spark写入
对于spark任务,建议开启AQE
spark.conf.set("spark.sql.adaptive.enabled","true")spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled","true")spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes","268435456")# 如果还产生大量小文件,在进行一次repartition# df.coalesce(5)# df.repartition($"ds") # 按字段分区输出,每个分区 1 个文件df.repartition(10)\.write \.partitionBy("dt")\.saveAsTable("table")repartition = 重新分区(可以增加 / 减少分区) + 全量 Shuffle + 数据均匀,成本较高。适合:增加并行度、需要均匀分区、大幅减少分区
coalesce:默认不 shuffle,只能减少分区,数据可能不均匀。适合:小幅减少分区、避免小文件、无 shuffle 需求
PySpark任务多种卡住问题
- 数据量不大,但某些task卡住几个小时,看Thread dump,卡在SocketInputStream
找到卡住的task对应的日志,udf有明显报错
pyspark worker已经挂掉,但是executor jvm一直在等待数据返回,导致卡住。 需要排查用户上游是否有脏数据,或者在udf中增加预期异常处理逻辑。 - 数据量大,一般是处理机器学习问题,数据量上亿级别,卡在某一个task。日志中有too large frame异常,一般是存在数据倾斜问题,某些task处理的数据量过大。
- Pyspark运行中报错内存超用 Current mem limits: xxx of max xxx
从pyspark的代码来看,python worker的内存使用并没有在executor中登记,也就python worker的内存使用是没有办法限制的,这就导致python worker的内存成为问题点。用户每个APP读取的数据量比较大,并且数据都通过python的UDF处理,因此有如下的日志:
这里可以看到python worker的内存开始有警告了,最终导致:memory ERROR,从而是整个qpp任务失败。
- 解决方案:
将用户的数据切成更小的文件,用多个app去处理,这样可以做到APP并行,同时每个APP处理的数据量比较小,并且可以完成整个任务。
- Python worker进程卡在读取shuffle数据
可以看到卡主的executor的堆栈在task 87102上读取socket数据,从日志中可以去看这个task的日志。
从日志中看到链接shuffle server后就没有了后面的日志,应该是卡在了读取shuffle数据上面,通过让运维排查shuffle节点状态,在处理问题。
- 解决方案:1)打开推断执行。2)排查shuffle节点后,重新运行任务。
- 没有名明显报错,单纯慢(像是卡住)
遇到这种问题,看用户的脚本是不是有问题,数据量是不是很大,用户的python udf函数是不是很多(用户python函数,还是注册的python udf),根据这些情况进行分情况处理。
最根本的处理要点就是,尽量用pyspark处理的数据量小。尽可能采用切分数据,app并行的方式跑pyspark任务。
PySpark其他问题
- PySpark写入偶发 Caused by: java.io.FileNotFoundException: File hdfs://xxx_xxx/0 does not exist
原因:通常是并发写入导致临时目录冲突
解决方法:请按照以下步骤尝试解决:
(1)检查是否存在多个任务同时并发写同一库表,避免同时写入
(2)若没有同时写入的情况,尝试设置参数 spark.sql.hive.convertMetastoreOrc=false ,重试任务
- PySpark报错Py4JJavaError: An error occurred while calling o205.jdbc. java.sql.SQLException: No suitable driver
原因:JDBC驱动加载失败
解决方法:请按照以下步骤尝试解决:
(1)如果是通过Toolkit访问PG外表,增加参数然后重试任务 spark.driver.defaultJavaOptions=-Djdbc.drivers=org.postgresql.Driver;spark.executor.defaultJavaOptions=-Djdbc.drivers=org.postgresql.Driver
(2)显式使用JDBC访问MySQL或者其他关系型数据库的表,以读取MySQL为例
defread_from_mysql(spark,query,table_ip,user_name,password,db_name):url='jdbc:mysql://{table_ip}:3306/{db_name}?useSSL=false'.format(table_ip=table_ip,db_name=db_name)table=query auth_mysql={"user":user_name,"password":password}data_df=spark.read.jdbc(url,table,properties=auth_mysql)returndata_df用户可以手动的指定JDBC的类型,例如:
在上面的auth_mysql中加入:“driver”:"com.mysql.jdbc.Driver,手动指定driver类型,也就是:auth_mysql = {“user”: user_name, “password”: password, “driver”:"com.mysql.jdbc.Driver}
这种修改的依据是:
Driver可以从用户指定的driver去获取
- PySpark的Python进程crash: Python worker exited unexpectedly (crashed)
原因:这种根据经验一般是Python进程由于占用内存太多被kill,Executor无法和Python进程通信导致。可在任务运行时由运维去物理机上进一步确认或者在Spark UI上通过Python dump和Cgroup确认:
物理机确认流程:cd /sys/fs/cgroup/memory/hadoop-yarn/container-xxx
在memory.stat文件看到Python进程因为OOM被kill
Spark UI在Executor页面通过Python dump和Cgroup确认
解决方法:请按照以下步骤尝试解决:
(1)通过增大分区数减少单个Python进程处理的数据量。尝试增大Shuffle分区数(默认200),调大Shuffle分区参数,例如spark.sql.shuffle.partitions=400 和 spark.default.parallelism=400。如果程序中使用coalesce或者repartition,可以尝试增大此方法的参数值来增加分区数
(2)调整Python进程可使用内存的参数(默认为1024M)spark.executor.memoryOverhead=4096,根据实际情况逐步增大
- PySpark报错Could not submit task to executor
原因:COS客户端通过线程池用来提交任务,当时线程池比较小时,导致提交任务被拒绝
从UI中打开用户卡主的Executor
从Executor的堆栈中定位到哪一个task卡主,然后从日志看卡主的日志的最后状态:
这里看到报错
排查代码后,发现COS客户端通过线程池用来提交任务,当时线程池比较小时,导致提交任务被拒绝
解决方法:通过以下参数设置cos线程池的大小:spark.hadoop.fs.cosn.upload_thread_pool=10
- PySpark PB数据解析报错TypeError: Descriptors cannot not be created directly
原因:Protobuf版本不兼容
解决方法:
针对有些需要使用PB来解析已经序列化写入的库表字段时,可以通过打印当前Python环境的PB版本来查看,然后使用对应的版本来生成PB协议文件,然后就能Python解析PB字符串了
importgoogle.protobufprint(google.protobuf.version)- PySpark的broadcast dump异常 Could not serialize broadcast: OverflowError: cannot serialize a string larger than
现象:File “…/pyspark.zip/pyspark/broadcast.py”, line 113, in dump pickle.dump(value, f, 2)
OverflowError: cannot serialize a string larger than 4GiB
OverflowError: cannot serialize a string larger than 4GiB
_pickle.PicklingError: Could not serialize broadcast: OverflowError: cannot serialize a string larger than 4GiB
原因:PySpark的broadcast依赖的pickle库的pickle.dump方法在pickling protocol过低的时候,不支持超过4G的对象的序列化
解决方法:在代码中添加如下代码,替换掉pyspark的broadcast.Broadcast.dump方法
frompysparkimportbroadcastimportpickledefbroadcast_dump(self,value,f):pickle.dump(value,f,4)# was 2, 4 is first protocol supporting >4GBf.close()returnf.name broadcast.Broadcast.dump=broadcast_dump参考链接:https://stackoverflow.com/questions/53371112/creating-parquet-petastorm-dataset-through-spark-fails-with-overflow-error-larg
- PySpark报错KeyError
File "/data11/yarnenv/local/usercache/hive/appcache/application_1231_123/container_e47_1231_123_01_000001/pyspark.zip/pyspark/rdd.py", line 1293, in takeUpToNumLeft File "WordCount.py", line 39, in <lambda> KeyError: u'15.xx\u7684\u9884\u5b9a\u4f1a\u8bae' File "WordCount.py", line 39, in <lambda> KeyError: u'15.xx\u7684\u9884\u5b9a\u4f1a\u8bae'原因:Key检索失败
解决方法:用户脚本抛出的运行异常,一般针对排查代码段能解决
- spark.executor.memoryOverhead 解释
详细解读 spark.executor.memoryOverhead 这个参数。它的逻辑与 Driver 的 Overhead 非常相似,但有一些针对 Executor 的特殊说明。核心定义:
- 参数名: spark.executor.memoryOverhead
- 核心含义: 为每个 Executor 进程分配的 额外内存 的大小,用于 JVM 堆外的开销。
详细解释
(1) 默认值如何计算?
executorMemory * spark.executor.memoryOverheadFactor, with minimum of spark.executor.minMemoryOverhead
- 计算公式: 与 Driver 端完全一致,是动态计算的。
- executorMemory:通过 --executor-memory 设置的 JVM 堆内存大小。
- spark.executor.memoryOverheadFactor:比例因子,默认也是 0.10(10%)。
- spark.executor.minMemoryOverhead:最小 Overhead 值,默认也是 384 MiB。
- 计算逻辑: 取 (executorMemory * overheadFactor) 和 minMemoryOverhead 中的 较大者。
(2)这个内存是用来做什么的?
Amount of additional memory… for things like VM overheads, interned strings, other native overheads, etc.
用途与 Driver 类似,用于 Executor 进程的 JVM 非堆开销:
- JVM 自身开销: 线程栈(每个运行的任务都会占用线程栈空间)、GC 数据结构、代码缓存等。
- 本地内存: Executor 可能使用的堆外缓冲区(例如,在进行 shuffle、排序或使用某些本地库时)。
(3) 这个开销的大小规律是什么?
This tends to grow with the executor size (typically 6-10%).
同样,Overhead 的大小与 Executor 的规模成正比。Executor 分配的内存越大、核心数越多(意味着线程越多),需要的 Overhead 也越大。
(4)在哪些集群模式下有效?
This option is currently supported on YARN and Kubernetes.
同样,只有在 YARN 或 Kubernetes 这类基于容器的集群管理器下,这个参数才至关重要,因为它直接关系到容器能否稳定运行而不被资源管理器“杀死”。
关键差异和重要说明(Note 部分)
Executor 的 Overhead 定义比 Driver 的更复杂,因为它明确包含了更多组件。Note 部分是全段的核心。
- 包含 PySpark Executor 的内存
Additional memory includes PySpark executor memory (when spark.executor.pyspark.memory is not configured)- 这是 Executor 与 Driver Overhead 的一个关键区别。
- 在 PySpark 应用中,每个 Executor 不仅有一个 JVM 进程,还有一个配套的 Python 进程(Python Worker) 来执行 Python 代码(例如 UDF)。
- 默认情况下,这个 Python 进程消耗的内存被计算在 memoryOverhead 之内。
- 只有当显式配置了 spark.executor.pyspark.memory 时,Python 进程的内存才会被单独管理,不再从 Overhead 中扣除。
- 包含同一容器内的其他非 Executor 进程
and memory used by other non-executor processes running in the same container.- 与 Driver 一样,容器内可能存在的其他辅助进程的内存也计入 Overhead。
- 容器总内存的最终计算公式(极其重要)
`The maximum memory size of container to running executor is determined by the sum of:- spark.executor.memoryOverhead
- spark.executor.memory
- spark.memory.offHeap.size
- spark.executor.pyspark.memory`
这是最关键的公式,它定义了向资源管理器申请的 Executor 容器总内存。它由四个部分相加组成:
容器总内存 = spark.executor.memory (JVM 堆内存)
- spark.memory.offHeap.size (Spark 管理的堆外内存,需手动开启)
- spark.executor.pyspark.memory (Python Worker 进程内存,如果配置了)
- spark.executor.memoryOverhead (其他所有额外内存)
重要关系图:
flowchatchart TD A[Executor Container Total Memory<br>向YARN/K8s申请的总内存] --> B[Spark JVM Heap<br>spark.executor.memory] A --> C[Managed Off-Heap<br>spark.memory.offHeap.size<br>(可选)] A --> D[Python Worker Memory<br>spark.executor.pyspark.memory<br>(可选,如不配置则计入Overhead)] A --> E[Memory Overhead<br>spark.executor.memoryOverhead<br>(包含JVM非堆/其他进程等)] D -.->|“如果不配置 (默认)”| E示意图解读: 容器总内存是四个部分之和。其中,Python工作进程内存是一个特殊部分:如果单独配置了,它独立存在;如果没配置,它就被包含在Memory Overhead里。
总结与配置建议
| 配置项 | 含义 | 默认值/示例 | 作用 |
|---|---|---|---|
| spark.executor.memory | JVM 堆内存 | 8g | 存储任务处理数据的 Java 对象 |
| spark.memory.offHeap.size | Spark 管理的堆外内存 | 0(默认关闭) | 存储序列化数据,减少 GC 压力 |
| spark.executor.pyspark.memory | Python 进程内存 | 未配置(默认计入 Overhead) | 单独控制 Python Worker 的内存 |
| spark.executor.memoryOverhead | 额外内存 | 自动计算(如 8g * 0.1 = 819MB) | 保障 JVM、Python 进程等稳定运行 |
| 容器总内存 | 实际向集群申请的内存 | 四者之和 | YARN/K8s 监控和限制的依据 |
何时需要手动调整 spark.executor.memoryOverhead?
- PySpark 应用(且未设置 spark.executor.pyspark.memory): 如果 Python UDF 处理大量数据,Python 进程会消耗巨量内存。你必须大幅提高 memoryOverhead 来避免容器被杀死。
- 出现内存溢出错误: 作业失败日志中出现 Container killed by YARN for exceeding memory limits,通常意味着 Overhead 不足,需要调高。
- Executor 负载很重: 如果 Executor 核心数多(线程多)、Shuffle 量大或使用了大量原生库,需要增加 Overhead。
- 启用堆外内存: 如果你设置了 spark.memory.offHeap.size=1g,理论上 Overhead 需要额外增加这 1GB。但根据 Note 中的公式,offHeap.size 是独立于 Overhead 的,所以你通常不需要为此调整 Overhead。但如果还有其他开销(如 Python),仍需增加。
最佳实践示例:
# 一个使用Python UDF的Spark应用,Executor配置示例spark-submit\--executor-memory 10g\--confspark.executor.memoryOverhead=3g\# 为Python进程和JVM开销预留充足内存--confspark.executor.pyspark.memory=2g\# 显式为Python进程分配2G,这2G不再从Overhead中扣除...核心要点: 理解 Executor 容器总内存的四个组成部分,并根据你的应用类型(纯 Scala/Java 还是 PySpark)和操作特点来合理分配这四部分的内存,是稳定运行 Spark 作业的关键。