PySpark原理介绍小文件处理背景hive 分区如果产生了大量小文件不仅会消耗存储元数据quota还会导致在读取该分区时性能和效率低下大量的时间浪费在了元数据获取上同时在数据存储上效率也偏低存储浪费在了元数据record上因此小文件是很值得进行优化的功能点。在没有shuffle阶段的处理过程中出现小文件通过distribute by cast(rand() * 10 as int) 增加一个shuffle阶段将小文件重分区成10个分区减少小文件。一.SQL写入1.单分区写入:作用全局随机打散均匀分区。适用场景解决数据倾斜但会产生大量小文件。单日数据 ≤ 2GB → distribute by cast(rand() as int) – 强制将所有数据给到 key0无论数据大小永远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.mapfilestrue;SEThive.merge.mapredfilestrue;SEThive.merge.size.per.task268435456;-- 256MBSEThive.merge.smallfiles.avgsize134217728;-- 128MBINSERTINTOTABLEtable1PARTITION(dt)SELECTcol1,col2,dtFROMtable2 DISTRIBUTEBYdt,rand()-- 按分区随机打散SORTBYdt,rand();-- 每个分区内有序避免碎片4.spark写入对于spark任务建议开启AQEspark.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.convertMetastoreOrcfalse 重试任务PySpark报错Py4JJavaError: An error occurred while calling o205.jdbc. java.sql.SQLException: No suitable driver原因JDBC驱动加载失败解决方法请按照以下步骤尝试解决1如果是通过Toolkit访问PG外表增加参数然后重试任务 spark.driver.defaultJavaOptions-Djdbc.driversorg.postgresql.Driverspark.executor.defaultJavaOptions-Djdbc.driversorg.postgresql.Driver2显式使用JDBC访问MySQL或者其他关系型数据库的表以读取MySQL为例defread_from_mysql(spark,query,table_ip,user_name,password,db_name):urljdbc:mysql://{table_ip}:3306/{db_name}?useSSLfalse.format(table_iptable_ip,db_namedb_name)tablequery auth_mysql{user:user_name,password:password}data_dfspark.read.jdbc(url,table,propertiesauth_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进程由于占用内存太多被killExecutor无法和Python进程通信导致。可在任务运行时由运维去物理机上进一步确认或者在Spark UI上通过Python dump和Cgroup确认物理机确认流程cd /sys/fs/cgroup/memory/hadoop-yarn/container-xxx在memory.stat文件看到Python进程因为OOM被killSpark UI在Executor页面通过Python dump和Cgroup确认解决方法请按照以下步骤尝试解决1通过增大分区数减少单个Python进程处理的数据量。尝试增大Shuffle分区数默认200调大Shuffle分区参数例如spark.sql.shuffle.partitions400 和 spark.default.parallelism400。如果程序中使用coalesce或者repartition可以尝试增大此方法的参数值来增加分区数2调整Python进程可使用内存的参数默认为1024Mspark.executor.memoryOverhead4096根据实际情况逐步增大PySpark报错Could not submit task to executor原因COS客户端通过线程池用来提交任务当时线程池比较小时导致提交任务被拒绝从UI中打开用户卡主的Executor从Executor的堆栈中定位到哪一个task卡主然后从日志看卡主的日志的最后状态这里看到报错排查代码后发现COS客户端通过线程池用来提交任务当时线程池比较小时导致提交任务被拒绝解决方法通过以下参数设置cos线程池的大小spark.hadoop.fs.cosn.upload_thread_pool10PySpark 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 4GiBOverflowError: 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.dumpbroadcast_dump参考链接https://stackoverflow.com/questions/53371112/creating-parquet-petastorm-dataset-through-spark-fails-with-overflow-error-largPySpark报错KeyErrorFile /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: u15.xx\u7684\u9884\u5b9a\u4f1a\u8bae File WordCount.py, line 39, in lambda KeyError: u15.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.1010%。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.memoryOverheadspark.executor.memoryspark.memory.offHeap.sizespark.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 Memorybr向YARN/K8s申请的总内存] -- B[Spark JVM Heapbrspark.executor.memory] A -- C[Managed Off-Heapbrspark.memory.offHeap.sizebr可选] A -- D[Python Worker Memorybrspark.executor.pyspark.memorybr可选如不配置则计入Overhead] A -- E[Memory Overheadbrspark.executor.memoryOverheadbr包含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.memoryOverheadPySpark 应用且未设置 spark.executor.pyspark.memory 如果 Python UDF 处理大量数据Python 进程会消耗巨量内存。你必须大幅提高 memoryOverhead 来避免容器被杀死。出现内存溢出错误 作业失败日志中出现 Container killed by YARN for exceeding memory limits通常意味着 Overhead 不足需要调高。Executor 负载很重 如果 Executor 核心数多线程多、Shuffle 量大或使用了大量原生库需要增加 Overhead。启用堆外内存 如果你设置了 spark.memory.offHeap.size1g理论上 Overhead 需要额外增加这 1GB。但根据 Note 中的公式offHeap.size 是独立于 Overhead 的所以你通常不需要为此调整 Overhead。但如果还有其他开销如 Python仍需增加。最佳实践示例# 一个使用Python UDF的Spark应用Executor配置示例spark-submit\--executor-memory 10g\--confspark.executor.memoryOverhead3g\# 为Python进程和JVM开销预留充足内存--confspark.executor.pyspark.memory2g\# 显式为Python进程分配2G这2G不再从Overhead中扣除...核心要点 理解 Executor 容器总内存的四个组成部分并根据你的应用类型纯 Scala/Java 还是 PySpark和操作特点来合理分配这四部分的内存是稳定运行 Spark 作业的关键。