PySpark多列合并为JSON列并修改JSON内字段名问题
解决PySpark聚合列转JSON时修改字段名的问题
方法1:聚合阶段直接指定目标字段名
在Spark SQL聚合查询时,直接将计算结果的别名设置为JSON需要的字段名,再通过to_json+struct打包成JSON列。
假设原数据表为business_data,包含date、cost、revenue、profit及分组字段(如product_id),示例SQL:
SELECT product_id, to_json( struct( SUM(CASE WHEN date >= add_months(current_date(), -3) THEN cost ELSE 0 END) AS cost_3m, SUM(CASE WHEN date >= add_months(current_date(), -3) THEN revenue ELSE 0 END) AS revenue_3m, SUM(CASE WHEN date >= add_months(current_date(), -3) THEN profit ELSE 0 END) AS profit_3m, SUM(CASE WHEN date >= add_months(current_date(), -6) THEN cost ELSE 0 END) AS cost_6m, SUM(CASE WHEN date >= add_months(current_date(), -6) THEN revenue ELSE 0 END) AS revenue_6m, SUM(CASE WHEN date >= add_months(current_date(), -6) THEN profit ELSE 0 END) AS profit_6m ) ) AS Last_Month_Details FROM business_data GROUP BY product_id
方法2:对已聚合的DataFrame重命名字段后打包
如果已经有包含last_3m_cost等6个聚合列的DataFrame,可直接在struct中为每个字段指定别名,无需额外创建中间列:
PySpark代码示例:
from pyspark.sql import functions as F # 假设df是已聚合后的DataFrame,包含目标6个列 df_final = df.withColumn( "Last_Month_Details", F.to_json( F.struct( F.col("last_3m_cost").alias("cost_3m"), F.col("last_3m_revenue").alias("revenue_3m"), F.col("last_3m_profit").alias("profit_3m"), F.col("last_6m_cost").alias("cost_6m"), F.col("last_6m_revenue").alias("revenue_6m"), F.col("last_6m_profit").alias("profit_6m") ) ) )
如果需要更结构化的JSON(比如将3个月和6个月数据分成子对象),可以嵌套struct:
df_final = df.withColumn( "Last_Month_Details", F.to_json( F.struct( F.struct( F.col("last_3m_cost").alias("cost"), F.col("last_3m_revenue").alias("revenue"), F.col("last_3m_profit").alias("profit") ).alias("last_3m"), F.struct( F.col("last_6m_cost").alias("cost"), F.col("last_6m_revenue").alias("revenue"), F.col("last_6m_profit").alias("profit") ).alias("last_6m") ) ) )
方法3:Spark SQL处理已聚合表
若已生成包含聚合列的临时表(如aggregated_data),可通过SQL直接重命名字段并打包:
SELECT product_id, to_json( struct( last_3m_cost AS cost_3m, last_3m_revenue AS revenue_3m, last_3m_profit AS profit_3m, last_6m_cost AS cost_6m, last_6m_revenue AS revenue_6m, last_6m_profit AS profit_6m ) ) AS Last_Month_Details FROM aggregated_data
内容的提问来源于stack exchange,提问作者user2258035
相关产品推荐
相关产品推荐

