如何在Databricks的PySpark中合并数据框并添加列?
在Databricks PySpark中合并表元数据并追加列的解决方案
核心结论
完全可以通过withColumn或表关联(Join)的方式为DataFrame追加列,以下是针对需求的具体实现步骤和代码:
步骤1:优化获取基础表元数据
先获取包含DatabaseName、TableName、IsTemporary的完整基础DataFrame,替代原有的重复循环逻辑:
from pyspark.sql.functions import lit from pyspark.sql.types import StructType, StructField, StringType from functools import reduce from pyspark.sql import DataFrame # 获取所有非临时表的元数据(如需包含临时表可去掉WHERE条件) base_df = spark.sql(""" SELECT database as DatabaseName, tableName as TableName, isTemporary as IsTemporary FROM spark_catalog.tables WHERE database IS NOT NULL """) base_df.createOrReplaceTempView("base_tables")
步骤2:关联DDL语句列
通过遍历基础表的记录,动态获取每个表的DDL,并通过withColumn添加关联字段(库名、表名),最后合并为统一的DDL DataFrame:
ddl_list = [] for row in base_df.collect(): db_name = row["DatabaseName"] table_name = row["TableName"] # 临时表无DDL,跳过处理 if not row["IsTemporary"]: # 获取单表DDL ddl_single = spark.sql(f"SHOW CREATE TABLE {db_name}.{table_name}") # 追加库名、表名字段用于关联 ddl_single = ddl_single.withColumn("DatabaseName", lit(db_name)) \ .withColumn("TableName", lit(table_name)) ddl_list.append(ddl_single) # 合并所有表的DDL数据 ddl_df = reduce(DataFrame.unionAll, ddl_list) if ddl_list else spark.createDataFrame([], schema=StructType([ StructField("crtmnt_stmnt", StringType()), StructField("DatabaseName", StringType()), StructField("TableName", StringType()) ])) ddl_df.createOrReplaceTempView("table_ddls")
步骤3:追加存储位置(Location)列
同样通过遍历获取每个表的存储位置,生成关联DataFrame后合并到主表:
location_list = [] for row in base_df.collect(): db_name = row["DatabaseName"] table_name = row["TableName"] if not row["IsTemporary"]: # 通过DESCRIBE EXTENDED获取表存储位置 desc_df = spark.sql(f"DESCRIBE EXTENDED {db_name}.{table_name}") location = desc_df.filter(desc_df.col_name == "Location").select("data_type").first()[0] # 生成包含关联字段的位置数据 loc_single = spark.createDataFrame([(db_name, table_name, location)], schema=["DatabaseName", "TableName", "Location"]) location_list.append(loc_single) location_df = reduce(DataFrame.unionAll, location_list) if location_list else spark.createDataFrame([], schema=StructType([ StructField("DatabaseName", StringType()), StructField("TableName", StringType()), StructField("Location", StringType()) ])) location_df.createOrReplaceTempView("table_locations")
步骤4:合并所有字段并保存
通过表关联将基础表、DDL表、位置表合并,得到最终结构的DataFrame,可直接保存为永久表:
# 方式1:用DataFrame API合并 final_df = base_df.join(ddl_df, on=["DatabaseName", "TableName"], how="left") \ .join(location_df, on=["DatabaseName", "TableName"], how="left") # 方式2:用Spark SQL合并 final_df = spark.sql(""" SELECT b.DatabaseName, b.TableName, b.IsTemporary, d.crtmnt_stmnt, l.Location FROM base_tables b LEFT JOIN table_ddls d ON b.DatabaseName = d.DatabaseName AND b.TableName = d.TableName LEFT JOIN table_locations l ON b.DatabaseName = l.DatabaseName AND b.TableName = l.TableName """) # 保存为永久表 final_df.write.mode("overwrite").saveAsTable("table_metadata_store")
关键说明
withColumn的核心作用是给DataFrame添加新列,这里用来给单表DDL数据追加库名、表名字段,确保后续能和基础表准确关联。- 原代码中循环创建临时表会覆盖之前的数据,导致无法保留所有表的元数据,优化后的逻辑通过收集所有表的结果再合并解决了这个问题。
- 临时表没有物理存储位置和DDL语句,所以处理时跳过了临时表的相关获取逻辑,避免报错。
内容的提问来源于stack exchange,提问作者Ria C
相关产品推荐
相关产品推荐

