PySpark如何通过循环按门店筛选生成多个独立DataFrame
PySpark按StoreName拆分多DataFrame实现方案
错误原因
你原有代码的语法错误是Python本身的限制:不支持在循环中直接使用df_{i}格式的动态字符串作为变量名赋值,和PySpark逻辑无关。
推荐实现方案
方案1:字典存储拆分结果(最常用,适配大数据量)
用字典存储所有拆分后的DataFrame,key为门店名称,value为对应门店的DataFrame,调用和批量处理都很方便:
from pyspark.sql.functions import col # 你原有提取门店列表的逻辑可以正常使用 store_list = df.select("StoreName").distinct().rdd.flatMap(lambda x:x).collect() store_df_map = {} for store in store_list: store_df_map[store] = df.where(col("StoreName") == store)
调用示例:获取ABC门店的DataFrame直接使用store_df_map["ABC"]即可。
方案2:直接分区落地存储(超大数据量优先)
如果你的最终目的是把不同门店的数据拆分存储,不需要在内存中持有多个DataFrame对象,直接用Spark原生分区写入即可,Spark内部会做计算优化,性能远高于手动循环筛选:
# 按StoreName分区写入为parquet格式,自动生成对应门店名称的子目录 df.write.partitionBy("StoreName").mode("overwrite").parquet("/your/target/path/")
方案3:动态注册全局变量(不推荐)
如果确实需要每个门店对应独立的变量名,可以通过globals()方法动态注册全局变量,执行后可直接调用df_ABC、df_PQR这类变量:
from pyspark.sql.functions import col store_list = df.select("StoreName").distinct().rdd.flatMap(lambda x:x).collect() for store in store_list: globals()[f"df_{store}"] = df.where(col("StoreName") == store)
注意:该写法维护性差,不利于后续批量迭代处理,非特殊需求不建议使用。
内容的提问来源于stack exchange,提问作者Bitanshu Das
相关产品推荐
相关产品推荐

