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

如何在Spark中基于Struct字段按sub-id聚合数据

解决方案

核心思路

要实现需求,关键是先将嵌套的Struct字段拆分为单行数据,按Id+Name+Car+sub-id分组收集同sub-id下的所有国家,再对US的记录补充其他国家列表,最后重新聚合回原行结构。


方案1:Spark SQL实现

步骤1:拆分嵌套Struct为单行

使用explode将每行的Struct数组拆成单独行:

SELECT Id, Name, Car, explode(Struct) AS struct_item
FROM original_table

步骤2:分组收集同sub-id的所有国家

按Id、Name、Car、sub-id分组,收集该组下的所有国家:

SELECT 
    Id, 
    Name, 
    Car, 
    struct_item.`sub-id` AS sub_id,
    collect_list(struct_item.country) AS all_countries
FROM (
    SELECT Id, Name, Car, explode(Struct) AS struct_item
    FROM original_table
) t
GROUP BY Id, Name, Car, struct_item.`sub-id`

步骤3:构造目标Struct数组

将拆分的行与分组结果关联,对US的记录添加other-country字段,最后聚合回原行:

SELECT 
    Id, 
    Name, 
    Car,
    array_agg(
        CASE 
            WHEN struct_item.country = 'US' THEN 
                struct(
                    struct_item.`sub-id` AS `sub-id`,
                    struct_item.country AS country,
                    array_remove(all_countries, 'US') AS `other-country`
                )
            ELSE struct_item
        END
    ) AS Struct
FROM (
    SELECT 
        t.Id, 
        t.Name, 
        t.Car,
        t.struct_item,
        g.all_countries
    FROM (
        SELECT Id, Name, Car, explode(Struct) AS struct_item
        FROM original_table
    ) t
    JOIN (
        SELECT 
            Id, 
            Name, 
            Car, 
            struct_item.`sub-id` AS sub_id,
            collect_list(struct_item.country) AS all_countries
        FROM (
            SELECT Id, Name, Car, explode(Struct) AS struct_item
            FROM original_table
        ) t
        GROUP BY Id, Name, Car, struct_item.`sub-id`
    ) g ON t.Id = g.Id AND t.Name = g.Name AND t.Car = g.Car AND t.struct_item.`sub-id` = g.sub_id
) final
GROUP BY Id, Name, Car

方案2:Pandas实现

步骤1:解析并拆分Struct字段

先将字符串格式的Struct转为字典列表,再拆分为单行:

import pandas as pd

# 解析Struct字符串为字典列表
def parse_struct(s):
    items = s.split('}, {')
    items = [item.strip('{}') for item in items]
    struct_list = []
    for item in items:
        kv_pairs = item.split(', ')
        sub_id = kv_pairs[0].split(': ')[1].strip('"')
        country = kv_pairs[1].split(': ')[1].strip('"')
        struct_list.append({'sub-id': sub_id, 'country': country})
    return struct_list

# 构造原始DataFrame
df = pd.DataFrame([
    [2, 'Alex', 'Porsche', '{sub-id: "car1", country: "US"}, {sub-id: "car1", country: "Mexico"}, {sub-id: "car-1", country: "Canada"}, {sub-id: "car-2", country: "Mexico"}'],
    [2, 'Alex', 'Toyota', '{sub-id: "car-2", country: "Germany"}, {sub-id: "car-3", country: "Brazil"}, {sub-id: "car-2", country: "US"}'],
    [3, 'Ben', 'BMW', '{sub-id: "car-2", country: "Germany"}, {sub-id: "car-3", country: "Brazil"}, {sub-id: "car-2", country: "US"}']
], columns=['Id', 'Name', 'Car', 'Struct'])

# 解析Struct并拆分为单行
df['Struct'] = df['Struct'].apply(parse_struct)
df_exploded = df.explode('Struct', ignore_index=True)

步骤2:分组收集同sub-id的国家列表

提取sub-id和country字段,分组收集该组下的所有国家:

# 提取sub-id和country到单独列
df_exploded['sub-id'] = df_exploded['Struct'].apply(lambda x: x['sub-id'])
df_exploded['country'] = df_exploded['Struct'].apply(lambda x: x['country'])

# 分组收集每个组的所有国家
country_groups = df_exploded.groupby(['Id', 'Name', 'Car', 'sub-id'])['country'].agg(list).reset_index(name='all_countries')

步骤3:构造目标Struct并聚合回原行

合并数据后,对US记录添加other-country,最后重新聚合:

# 合并原数据与分组结果
df_merged = pd.merge(df_exploded, country_groups, on=['Id', 'Name', 'Car', 'sub-id'])

# 构造新的Struct
def build_new_struct(row):
    struct = row['Struct'].copy()
    if struct['country'] == 'US':
        other_countries = [c for c in row['all_countries'] if c != 'US']
        struct['other-country'] = other_countries
    return struct

df_merged['new_Struct'] = df_merged.apply(build_new_struct, axis=1)

# 聚合回原行结构
result_df = df_merged.groupby(['Id', 'Name', 'Car'])['new_Struct'].agg(list).reset_index(name='Struct')

内容的提问来源于stack exchange,提问作者the_stack_over_flew

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 16:12:02