news 2026/8/22 13:39:08

PySpark小文件处理+卡点问题

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
PySpark小文件处理+卡点问题

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任务多种卡住问题

  1. 数据量不大,但某些task卡住几个小时,看Thread dump,卡在SocketInputStream

    找到卡住的task对应的日志,udf有明显报错

    pyspark worker已经挂掉,但是executor jvm一直在等待数据返回,导致卡住。 需要排查用户上游是否有脏数据,或者在udf中增加预期异常处理逻辑。
  2. 数据量大,一般是处理机器学习问题,数据量上亿级别,卡在某一个task。日志中有too large frame异常,一般是存在数据倾斜问题,某些task处理的数据量过大。
  3. 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处理的数据量比较小,并且可以完成整个任务。
  1. Python worker进程卡在读取shuffle数据

    可以看到卡主的executor的堆栈在task 87102上读取socket数据,从日志中可以去看这个task的日志。

    从日志中看到链接shuffle server后就没有了后面的日志,应该是卡在了读取shuffle数据上面,通过让运维排查shuffle节点状态,在处理问题。
  • 解决方案:1)打开推断执行。2)排查shuffle节点后,重新运行任务。
  1. 没有名明显报错,单纯慢(像是卡住)
    遇到这种问题,看用户的脚本是不是有问题,数据量是不是很大,用户的python udf函数是不是很多(用户python函数,还是注册的python udf),根据这些情况进行分情况处理。
    最根本的处理要点就是,尽量用pyspark处理的数据量小。尽可能采用切分数据,app并行的方式跑pyspark任务。

PySpark其他问题

  1. PySpark写入偶发 Caused by: java.io.FileNotFoundException: File hdfs://xxx_xxx/0 does not exist

原因:通常是并发写入导致临时目录冲突
解决方法:请按照以下步骤尝试解决:
(1)检查是否存在多个任务同时并发写同一库表,避免同时写入
(2)若没有同时写入的情况,尝试设置参数 spark.sql.hive.convertMetastoreOrc=false ,重试任务

  1. 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去获取

  1. 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,根据实际情况逐步增大

  1. PySpark报错Could not submit task to executor

原因:COS客户端通过线程池用来提交任务,当时线程池比较小时,导致提交任务被拒绝
从UI中打开用户卡主的Executor

从Executor的堆栈中定位到哪一个task卡主,然后从日志看卡主的日志的最后状态:
这里看到报错

排查代码后,发现COS客户端通过线程池用来提交任务,当时线程池比较小时,导致提交任务被拒绝
解决方法:通过以下参数设置cos线程池的大小:spark.hadoop.fs.cosn.upload_thread_pool=10

  1. PySpark PB数据解析报错TypeError: Descriptors cannot not be created directly

原因:Protobuf版本不兼容
解决方法:
针对有些需要使用PB来解析已经序列化写入的库表字段时,可以通过打印当前Python环境的PB版本来查看,然后使用对应的版本来生成PB协议文件,然后就能Python解析PB字符串了

importgoogle.protobufprint(google.protobuf.version)
  1. 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

  1. 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检索失败
解决方法:用户脚本抛出的运行异常,一般针对排查代码段能解决

  1. 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 部分是全段的核心。

  1. 包含 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 中扣除。
  2. 包含同一容器内的其他非 Executor 进程
    and memory used by other non-executor processes running in the same container.
    • 与 Driver 一样,容器内可能存在的其他辅助进程的内存也计入 Overhead。
  3. 容器总内存的最终计算公式(极其重要)
    `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.memoryJVM 堆内存8g存储任务处理数据的 Java 对象
spark.memory.offHeap.sizeSpark 管理的堆外内存0(默认关闭)存储序列化数据,减少 GC 压力
spark.executor.pyspark.memoryPython 进程内存未配置(默认计入 Overhead)单独控制 Python Worker 的内存
spark.executor.memoryOverhead额外内存自动计算(如 8g * 0.1 = 819MB)保障 JVM、Python 进程等稳定运行
容器总内存实际向集群申请的内存四者之和YARN/K8s 监控和限制的依据

何时需要手动调整 spark.executor.memoryOverhead?

  1. PySpark 应用(且未设置 spark.executor.pyspark.memory): 如果 Python UDF 处理大量数据,Python 进程会消耗巨量内存。你必须大幅提高 memoryOverhead 来避免容器被杀死。
  2. 出现内存溢出错误: 作业失败日志中出现 Container killed by YARN for exceeding memory limits,通常意味着 Overhead 不足,需要调高。
  3. Executor 负载很重: 如果 Executor 核心数多(线程多)、Shuffle 量大或使用了大量原生库,需要增加 Overhead。
  4. 启用堆外内存: 如果你设置了 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 作业的关键。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/22 13:34:39

X-AnyLabeling 数据标注工具快速上手指南:AI 预标注 + 手动修正

X-AnyLabeling 数据标注工具快速上手指南&#xff1a;AI 预标注 手动修正 【免费下载链接】X-AnyLabeling X-AnyLabeling: A lightweight, efficient, and unified cross-platform desktop application for annotating text, image, video, and multimodal data, combining ve…

作者头像 李华
网站建设 2026/8/22 13:32:30

C++的Concept/Model模式

在 C++(尤其是 C++20 引入 Concepts 之后)中,Concept/Model(概念/模型)是一种非常强大的泛型编程范式。它不仅能让你的模板代码更安全、可读性更高,还能实现优雅的运行时多态(类似于接口,但没有虚函数表的开销)。 为了让你、、彻底搞懂它,我们把它拆成两部分来看: …

作者头像 李华
网站建设 2026/8/22 13:28:38

SMU调试工具:AMD Ryzen 处理器底层调试与核心调参完整指南

SMU调试工具&#xff1a;AMD Ryzen 处理器底层调试与核心调参完整指南 【免费下载链接】SMUDebugTool A dedicated tool to help write/read various parameters of Ryzen-based systems, such as manual overclock, SMU, PCI, CPUID, MSR and Power Table. 项目地址: https:…

作者头像 李华
网站建设 2026/8/22 13:27:37

30秒跑通 ResNet-50:本地图片识别 1000 类的最短路径

30秒跑通 ResNet-50&#xff1a;本地图片识别 1000 类的最短路径 【免费下载链接】resnet-50 项目地址: https://ai.gitcode.com/hf_mirrors/microsoft/resnet-50 想给照片自动打标&#xff1f;ResNet-50 图像分类模型能在你本地完成 1000 类 ImageNet 识别&#xff0c…

作者头像 李华
网站建设 2026/8/22 13:27:00

AI Agent运行时安全:构建基于LLM-as-Judge的动态评估与拦截框架

1. 项目概述&#xff1a;当AI Agent开始“自作主张”&#xff0c;我们如何确保安全&#xff1f;最近和几个做AI应用落地的朋友聊天&#xff0c;大家不约而同地提到了同一个痛点&#xff1a;AI Agent&#xff08;智能体&#xff09;在调用外部工具&#xff08;Tool Use&#xff…

作者头像 李华