You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 03:33:22