Spark字典存储的DataFrame循环新增列不生效问题及思路验证
问题原因与方案解答
1. 字典中新增列丢失的原因
Spark DataFrame是不可变数据结构,withColumn操作不会修改原DF,而是返回一个包含新列的全新DF。你写的循环里,x只是遍历df_dict.values()时生成的临时变量,给x赋值只是修改了临时变量的指向,完全不会更新字典中存储的原DF对象。
想要把修改后的DF存回字典,按如下方式修改循环即可:
for key in df_dict: df_dict[key] = df_dict[key].withColumn('new_col', udf_some_func(df_dict[key]['col1'], df_dict[key]['col2']))
2. 先拆分DF到字典再处理的思路问题
该思路效率反而更低,不推荐使用:
- 把大DF拆成9个小DF存到字典,会额外增加Spark的调度开销,多个小DF的计算任务无法被Spark合并优化,实际运行速度比全量DF统一处理更慢
- 按
naics_code输出不同CSV的需求,Spark原生write.partitionBy接口完全支持,不需要手动拆分DF
推荐的最优实现方式:
# 先对全量DF统一新增需要的列 df_with_new_col = df.withColumn('new_col', udf_some_func(df['col1'], df['col2'])) # 直接按naics_code分区输出,会自动生成每个naics_code对应的子文件夹,里面是该分区的CSV文件 df_with_new_col.write\ .option("header", True)\ .partitionBy("naics_code")\ .mode("overwrite")\ .csv("输出目录路径")
如果要求每个naics_code对应单个CSV文件,新增repartition操作即可:
df_with_new_col.repartition("naics_code")\ .write\ .option("header", True)\ .partitionBy("naics_code")\ .mode("overwrite")\ .csv("输出目录路径")
内容的提问来源于stack exchange,提问作者9ganzi
相关产品推荐
相关产品推荐

