Spark2.4.5下如何动态将JSON宽表转换为(id,name)窄表?
嘿,这个场景我太熟悉了——硬编码stack()的参数不仅麻烦,遇到列数变化还得改代码,完全不灵活。给你两个靠谱的动态实现方案,适配Spark 2.4.5:
方案一:动态生成SQL的
stack()语句 这个思路是利用Spark的元数据自动获取列名,再拼接出stack()需要的参数,完全不用手动写每一列:
获取表的所有列名
先通过show columns语句拿到宽表的列名列表:# 假设你的宽表已经注册为临时表`table_name` column_rows = spark.sql("show columns from table_name").collect() column_names = [row["col_name"] for row in column_rows]自动拼接
stack()的参数
每个列需要对应'列名',列名``的格式,我们把所有列的这个组合拼接起来:stack_items = [f"'{col}', `{col}`" for col in column_names] stack_clause = f"stack({len(stack_items)}, {','.join(stack_items)}) as (id, name)"执行动态生成的SQL
直接用拼接好的语句查询就能得到窄表:narrow_df = spark.sql(f"SELECT {stack_clause} FROM table_name") narrow_df.show()
方案二:用DataFrame API实现(更推荐)
如果你习惯用Spark的DataFrame链式操作,这个方法更简洁,不需要拼接SQL字符串,也不用依赖临时表:
核心思路是把宽表的所有列转换成一个Map类型的列,再用explode把Map拆分成键值对(也就是我们要的id和name):
from pyspark.sql.functions import create_map, explode, lit from itertools import chain # 假设`df`是你加载后的宽表Dataset/DataFrame # 生成Map表达式:键是列名(字符串常量),值是列的实际内容 map_expr = create_map(list(chain(*[(lit(col), df[col]) for col in df.columns]))) # 拆分Map得到窄表 narrow_df = df.select(explode(map_expr).alias("id", "name")) narrow_df.show()
解释一下:create_map需要成对的键和值,chain(*...)是把嵌套的列表摊平成一维;lit(col)把列名转成字符串,和列本身组成Map的键值对;explode会把Map的每个键值对拆成单独一行,刚好就是目标的窄表结构。
两种方案都能实现动态转换,不用硬编码任何列名,完全满足你的需求~
内容的提问来源于stack exchange,提问作者David
相关产品推荐
相关产品推荐

