如何无需for循环,基于基础表列值拆分生成多个子表?
无需循环拆分Spark表为多个子表的方法
针对你的需求,确实有不用for循环的解决方案,以下两种方法可以实现:
方法一:生成并批量执行SQL建表语句
通过Spark SQL的聚合函数一次性生成所有创建子表的SQL语句,再批量执行,全程无需循环:
# 生成所有创建子表的SQL语句 create_table_commands = spark.sql(""" SELECT concat_ws('; ', collect_list( concat( "CREATE TABLE IF NOT EXISTS ", Category, " AS SELECT * FROM df_transportation WHERE Category='", Category, "'" ) )) AS sql_commands FROM (SELECT DISTINCT Category FROM df_transportation) """).first()["sql_commands"] # 批量执行SQL语句,生成所有子表 spark.sql(create_table_commands)
方法二:使用Delta分区表(更适合大数据场景)
如果可以接受以分区形式存储而非独立表,这种方法性能更优,且无需循环:
# 开启动态分区配置(确保分区覆盖模式正确) spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic") spark.conf.set("hive.exec.dynamic.partition", "true") spark.conf.set("hive.exec.dynamic.partition.mode", "nonstrict") # 将原表按Category分区保存为Delta表 df_transportation.write.mode("overwrite")\ .format("delta")\ .partitionBy("Category")\ .saveAsTable("transportation_partitioned")
后续查询指定类别数据时,只需执行SELECT * FROM transportation_partitioned WHERE Category='car'即可,Spark会自动读取对应分区的数据,效率远高于拆分独立表。
注:你提供的原循环代码存在两处错误:一是
col("category") == "category"中误将变量写成字符串常量,二是saveAsTable方法的括号未闭合,修正后代码如下:
df_tmp = df_transportation.groupBy("category").count() categories = list(df_tmp.select("category").toPandas()["category"]) for category in categories: df_cat = df_transportation.filter(col("category") == category) df_cat.write.mode("overwrite").format("delta").saveAsTable(f'{category}')
内容的提问来源于stack exchange,提问作者siva
相关产品推荐
相关产品推荐

