PySpark DataFrame操作卡顿求助:过滤/转Pandas/写入CSV异常
PySpark DataFrame 写入/特定过滤操作停滞问题排查与解决思路
核心问题梳理
12万条记录的PySpark DataFrame,执行等值过滤(如city == 'Hannover, Landeshauptstadt')正常,但执行写入CSV、转换为Pandas DF、两列不等值过滤时耗时极长甚至停滞;调整集群配置后,limit(4)写入成功,但limit(500)仍停滞,show()操作正常。
排查与解决步骤
一、检查数据本身的异常
验证列数据类型与空值
两列不等值过滤会触发全表扫描,若列类型不一致(如state为字符串、state_osm为数值)或存在大量空值,会额外增加计算开销:# 查看列类型 df_h_with_distance_order_osm_P.select("state", "state_osm").printSchema() # 统计空值数量 from pyspark.sql.functions import col df_h_with_distance_order_osm_P.filter(col("state").isNull() | col("state_osm").isNull()).count()若存在类型不匹配,先统一类型(如
cast(StringType()));若空值过多,可先过滤空值再执行目标操作。分析执行计划
复杂的DataLineage可能导致Spark优化器无法生成高效执行计划,用explain查看瓶颈:df_h_with_distance_order_osm_P.explain(True)重点关注是否存在不必要的Shuffle操作、全表扫描的重复执行,或谓词下推失效的情况。
二、针对写入操作的优化
调整写入分区数
即使是limit(500),若原DF分区数过多,写入时会生成大量小文件,触发IO瓶颈。用coalesce(窄依赖,无Shuffle)合并分区后再写入:df_h_with_distance_order_osm_P.limit(500).coalesce(1).write.csv("/FileStore/tables/final_s.csv", header=True)排查存储路径与权限
测试写入临时路径(如/tmp/test_output),排除目标路径的权限限制、存储系统吞吐量瓶颈或文件冲突问题。临时关闭动态分配
动态分配在小数据量场景下可能导致Executor启动/销毁的额外开销,临时关闭该参数后测试:spark.conf.set("spark.dynamicAllocation.enabled", "false")
三、优化过滤操作的执行效率
- 收集表统计信息
等值过滤正常可能是因为city列有有效统计信息,Spark能做谓词下推;而两列不等过滤无法利用该优化,需手动收集统计信息帮助优化器:
之后用SQL语法执行过滤,可能获得更优执行计划:df_h_with_distance_order_osm_P.createOrReplaceTempView("df_view") spark.sql("ANALYZE TABLE df_view COMPUTE STATISTICS FOR ALL COLUMNS")state_not_equal = spark.sql("SELECT * FROM df_view WHERE state != state_osm")
四、修正缓存/持久化的用法
- 缓存需触发计算才会生效,仅调用
cache()不会实际缓存数据,需配合count()或show()触发:
若内存不足,改用df_cached = df_h_with_distance_order_osm_P.cache() df_cached.count() # 触发缓存MEMORY_AND_DISK持久化级别:from pyspark.storagelevel import StorageLevel df_cached.persist(StorageLevel.MEMORY_AND_DISK)
五、利用Spark UI定位瓶颈
在Databricks或Spark UI中查看以下信息:
- Jobs/Stages:确认停滞阶段是Shuffle、数据读取还是写入环节
- Tasks:查看是否有任务长时间未完成(如数据倾斜导致部分Task处理远超平均的数据量)
- Executor Metrics:检查内存、CPU使用率,确认是否存在资源不足
内容的提问来源于stack exchange,提问作者Mee
相关产品推荐
相关产品推荐

