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

Spark-SQL未知Schema下实现行列转置及异Schema表合并

用Spark-SQL解析不规则JSON脏数据并规整化

原始数据

{
   "user_1":{
      "user_name":"Andy",
      "user_lastname":"Martins",
      "user_accounts":{
         "facebook":"https://...",
         "twitter":[
            "https://...",
            "https://..."
         ]
      }
   },
   "user_2":{
      "user_name":"Joe",
      "user_lastname":"Foo",
      "user_accounts":{
         "twitter":"https://",
         "stack_overflow":"https://..."
      }
   }
}

Spark读取后的初始结果

val df = spark
    .read
    .option("multiline", "true")
    .option("inferSchema", "false")
    .format("json").load("data.json")

df.show(false)
// +----------------------------------------------------------+-----------------------------------+
// |user_1                                                    |user_2                             |
// +----------------------------------------------------------+-----------------------------------+
// |{{https://..., [https://..., https://...]}, Martins, Andy}|{{https://..., https://}, Foo, Joe}|
// +----------------------------------------------------------+-----------------------------------+

期望规整结果

+-------+---------+-------------+-------------+--------------------------+-------------+
|user_id|user_name|user_lastname|fb_account   |tw_account                |so_account   |
+-------+---------+-------------+-------------+--------------------------+-------------+
|user_1 |Andy     |Martins      |"https://..."|["https://...","https://..|             |
|user_2 |Joe      |Foo          |             |"https://..."             |"https://..."|
+-------+---------+-------------+-------------+--------------------------+-------------+

核心问题

  1. 百万级列的转置问题:直接获取所有列名并循环处理会导致内存溢出,无法高效完成列转行。
  2. Schema不一致无法合并:转置后每个用户的data结构体Schema不同,无法用UNION ALL合并,强制转字符串会丢失数组类型。

解决方案

1. 高效列转行(避免百万级列内存问题)

无需遍历每个列生成临时表,直接用stack函数动态生成转换逻辑,一次性完成所有列的转置:

val columns = df.columns
// 生成stack表达式:stack(N, '列名1', 列1, '列名2', 列2, ...)
val stackExpr = s"stack(${columns.length}, ${columns.flatMap(c => s"'$c', `$c`").mkString(", ")}) as (user_id, data)"
val pivotedDf = df.selectExpr(stackExpr)

这种方式避免了循环处理单个列的性能开销,适合百万级列的场景。

2. 统一Schema并保留原始类型

通过to_json将不一致的结构体转为JSON字符串,再用预定义的统一Schema解析,同时处理字段类型不一致的情况:

import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._

// 定义统一的目标Schema,兼容所有可能的字段
val userSchema = StructType(Seq(
  StructField("user_name", StringType),
  StructField("user_lastname", StringType),
  StructField("user_accounts", StructType(Seq(
    StructField("facebook", StringType),
    StructField("twitter", ArrayType(StringType)), // 统一为数组类型
    StructField("stack_overflow", StringType)
  )))
))

// 解析并规整字段
val finalDf = pivotedDf
  .select(
    $"user_id",
    from_json(to_json($"data"), userSchema).alias("user_data")
  )
  .select(
    $"user_id",
    $"user_data.user_name",
    $"user_data.user_lastname",
    $"user_data.user_accounts.facebook".alias("fb_account"),
    // 将单个字符串的twitter转为数组,保证类型一致
    when(
      size($"user_data.user_accounts.twitter") === 0,
      array($"user_data.user_accounts.twitter")
    ).otherwise($"user_data.user_accounts.twitter").alias("tw_account"),
    $"user_data.user_accounts.stack_overflow".alias("so_account")
  )

finalDf.show(false)
  • to_json + from_json可以自动填充缺失字段为null,解决Schema不一致问题
  • 用when处理twitter字段的类型差异,确保所有行的字段类型统一

3. dbt/Jinja SQL实现

如果用dbt处理,通过Jinja宏动态生成stack参数,避免手动编写大量代码:

{% set source_table = ref('jsons') %}
{% set columns = adapter.get_columns_in_relation(source_table) %}
{% set stack_params = [] %}
{% for col in columns %}
  {% do stack_params.append("'" ~ col.column ~ "'") %}
  {% do stack_params.append("`" ~ col.column ~ "`") %}
{% endfor %}

-- 列转行
CREATE OR REPLACE TEMP VIEW pivoted_data AS
SELECT stack({{ columns|length }}, {{ stack_params|join(', ') }}) AS (user_id, data)
FROM {{ source_table }};

-- 解析并规整
CREATE OR REPLACE TEMP VIEW final_data AS
SELECT
  user_id,
  user_data.user_name,
  user_data.user_lastname,
  user_data.user_accounts.facebook AS fb_account,
  CASE
    WHEN size(user_data.user_accounts.twitter) = 0 THEN array(user_data.user_accounts.twitter)
    ELSE user_data.user_accounts.twitter
  END AS tw_account,
  user_data.user_accounts.stack_overflow AS so_account
FROM (
  SELECT
    user_id,
    from_json(
      to_json(data),
      'struct<user_name:string,user_lastname:string,user_accounts:struct<facebook:string,twitter:array<string>,stack_overflow:string>>'
    ) AS user_data
  FROM pivoted_data
);

内容的提问来源于stack exchange,提问作者Eric Ávila

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 11:42:24