Spark中如何按列前4字符分组,同步聚合描述列并保留Null?
问题描述
我正在根据列的前4个字符对数据进行分类,前4字符相同的归为一类。当前使用以下Spark代码实现:
# 提取标题的前4个元素 df1_us2 = df1_us2.withColumn("first_2_char", df1_us2.clean_company_name.substr(1,4)) # 分组并聚合为列表 group_user = df1_us2.groupBy('first_2_char').agg(collect_set('col1').alias('col11'))
每条数据的col1对应一个description列,我希望对description列也执行相同的分类聚合,要求聚合后的两个列表长度必须一致(即使description为Null也要保留),通过索引匹配对应描述。
输入示例:
| col1 | description |
|---|---|
| summer | a season |
| summary | it is a brief |
| common | having similar |
| communication | null |
| house | living place |
期望输出:
| col11 | description1 |
|---|---|
| ['summer','summary'] | ['a season','it is a brief'] |
| ['common','communication'] | ['having similar', null] |
| ['house'] | ['living place'] |
解决方案
需要修改聚合逻辑,不能使用collect_set——它会去重并打乱顺序,导致两个列表长度不匹配。应使用collect_list保留原始顺序和所有元素(包括Null值),同时对col1和description分别聚合:
from pyspark.sql.functions import collect_list, substr # 提取前4个字符(修正变量名避免混淆) df1_us2 = df1_us2.withColumn("first_4_char", df1_us2.clean_company_name.substr(1, 4)) # 分组后同时聚合col1和description为对应列表 group_user = df1_us2.groupBy('first_4_char').agg( collect_list('col1').alias('col11'), collect_list('description').alias('description1') )
关键说明
- 替换
collect_set为collect_list:collect_list会按数据原始顺序收集所有元素,包括Null,确保col11和description1的长度完全一致,索引一一对应。 - 在
agg方法中新增对description的聚合操作,直接生成目标列description1。 - 原代码中分组列变量名
first_2_char与实际取前4字符的逻辑不符,建议改为first_4_char提升可读性。
内容的提问来源于stack exchange,提问作者user15649753
相关产品推荐
相关产品推荐

