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://..."| +-------+---------+-------------+-------------+--------------------------+-------------+
核心问题
- 百万级列的转置问题:直接获取所有列名并循环处理会导致内存溢出,无法高效完成列转行。
- 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
相关产品推荐
相关产品推荐

