Databricks使用PySpark读取列乱序CSV重排列报RLock错误
错误根因
触发cannot pickle '_thread.RLock' object报错有两个直接原因:
- 错误使用
spark.sparkContext.parallelize(claimdenials_df_raw):parallelize方法仅支持传入本地Python可序列化集合,不能直接传入DataFrame对象。DataFrame本身绑定了SparkContext的线程锁资源,强行传入做序列化分发时就会触发锁对象无法序列化的错误。 - 代码存在笔误:后续调用
rdd2=df.rdd.map时的df变量从未被定义。
另外你完全不需要转RDD实现列重排:Spark读取带表头的CSV时,会自动按列名做映射,和源文件里列的物理顺序无关,用DataFrame原生API就可以直接完成固定列序的输出,性能远高于RDD map实现。
正确实现代码
直接通过select方法按目标顺序传入列名即可,无论源CSV列是C/A/B还是其他任意乱序,都会自动对齐列名输出:
from pyspark.sql import SparkSession from pyspark.sql.functions import * appName = "ColumnReorder" spark = SparkSession.Builder().appName(appName).getOrCreate() parentname='xx' filename='Test (2)' extension='.csv' filePath=f"dbfs:/mnt/bronze/landing/x/{parentname}/current/{filename}{extension}" # 读取CSV文件 claimdenials_df_raw = (spark .read .format("csv") .option("multiLine", "true") .option("header", "true") .option("escape", '"') .load(filePath)) # 定义你需要的固定列顺序,注意列名必须和CSV表头完全匹配(包括前后空格) target_column_order = [ "Id", "First Name", "Last Name", "Date of Birth", "PatientId" ] # 选列完成重排 df_correct_order = claimdenials_df_raw.select(*target_column_order) display(df_correct_order)
注意事项
- 列名必须和CSV表头完全一致,你原代码里
" Last Name"列名前多了一个前置空格,要和实际表头字符对齐,否则会报列不存在错误。 - 不要在常规数据处理场景下随意把DataFrame转成RDD:RDD操作会跳过Spark的Catalyst优化器,数据需要在JVM和Python进程之间做序列化拷贝,数据量较大时性能会出现量级下降,还容易触发各类序列化问题。
内容的提问来源于stack exchange,提问作者Vaishnavi S
相关产品推荐
相关产品推荐

