如何在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
相关产品推荐
相关产品推荐

