Spark DataFrame通用函数实现:列转行并生成字典映射
通用Spark DataFrame列转行生成映射字典函数
需求
将Spark DataFrame中指定列转为行,以gw列的值作为字典的键,其他指定列的对应值作为字典的值,需要实现一个可接收多列输入的通用函数。
原始DataFrame
+----------------+-------+-------+ | gw | rrc | re_est| +----------------+-------+-------+ |210.142.27.137 |1400.0 |26.0 | |210.142.27.202 |2300 |12 | +----------------+-------+-------+
期望输出
+-------+-------------------------------------------------------+ | index | gw_mapping | +-------+-------------------------------------------------------+ | rrc |{210.142.27.137:1400.0, 210.142.27.202:2300} | | re_est|{210.142.27.137:26.0, 210.142.27.202:12} | +-------+-------------------------------------------------------+
已尝试的硬编码实现
import pyspark.sql.functions as F result_df = ( df .select('gw', F.expr("stack(2, 'rrc', rrc, 're_est', re_est) AS (index, value)")) .groupby('index') .agg(F.expr("map_from_entries(collect_list(struct(gw, value))) as gw_mapping")) )
通用函数实现
原代码的问题在于stack表达式是硬编码的,无法适配任意数量的输入列。下面是动态生成stack表达式的通用函数:
import pyspark.sql.functions as F def cols_to_gw_mapping(df, gw_col, target_cols): # 动态生成stack的参数:每个列对应('列名', 列)的组合 stack_items = [] for col in target_cols: stack_items.append(f"'{col}'") stack_items.append(col) stack_expr = f"stack({len(target_cols)}, {', '.join(stack_items)}) AS (index, value)" return ( df .select(gw_col, F.expr(stack_expr)) .groupBy('index') .agg(F.map_from_entries(F.collect_list(F.struct(gw_col, 'value'))).alias('gw_mapping')) )
使用示例
# 假设原始DataFrame为df result_df = cols_to_gw_mapping(df, gw_col='gw', target_cols=['rrc', 're_est']) result_df.show(truncate=False)
说明
- 函数接收三个参数:原始DataFrame、作为字典键的列名
gw_col、需要转换的目标列列表target_cols - 动态构造
stack表达式:根据目标列的数量生成对应的参数,避免硬编码 - 使用
map_from_entries结合collect_list(struct(gw_col, value))生成每个index对应的gw映射字典
内容的提问来源于stack exchange,提问作者sam
相关产品推荐
相关产品推荐

