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

如何将Spark DataFrame聚合为含键值对的单列(支持字符串/列表/映射)

实现Spark DataFrame转单键值对集合列的方案

当然可以用Spark DataFrame的内置函数实现,完全不需要自定义UDF!这里给你几种对应不同格式要求的实现方案,按需选择就行:

1. 转换为Map类型(推荐,原生键值对结构)

这种方式生成的Aggregated列是Map[String, Int]类型,直接对应键值对集合的语义,后续处理也更方便:

import org.apache.spark.sql.functions._

// 假设原DataFrame名为originalDF
val resultDF = originalDF
  // 先将protocol转为小写
  .withColumn("protocol_lower", lower(col("protocol")))
  // 聚合生成Map类型列
  .agg(
    map_from_entries(collect_list(struct(col("protocol_lower"), col("count")))).alias("Aggregated")
  )

执行后resultDF的结构如下:

+------------------------------------------------------------------+
|Aggregated                                                         |
+------------------------------------------------------------------+
|{tcp -> 8231, icmp -> 7314, udp -> 5523, igmp -> 4423, egp -> 2331}|
+------------------------------------------------------------------+

2. 转换为键值对字符串列表

如果需要列表格式,每个元素是protocol: count格式的字符串:

import org.apache.spark.sql.functions._

val resultDF = originalDF
  // 拼接单个键值对字符串
  .withColumn("kv_pair", concat(lower(col("protocol")), lit(": "), col("count")))
  // 收集为列表
  .agg(collect_list(col("kv_pair")).alias("Aggregated"))

生成的结果结构:

+-----------------------------------------------------------+
|Aggregated                                                 |
+-----------------------------------------------------------+
|[tcp: 8231, icmp: 7314, udp: 5523, igmp: 4423, egp: 2331]|
+-----------------------------------------------------------+

3. 转换为单个拼接字符串

如果需要把所有键值对合并成一个字符串:

import org.apache.spark.sql.functions._

val resultDF = originalDF
  .withColumn("kv_pair", concat(lower(col("protocol")), lit(": "), col("count")))
  // 用逗号分隔拼接所有键值对
  .agg(concat_ws(", ", collect_list(col("kv_pair"))).alias("Aggregated"))

生成的结果结构:

+---------------------------------------------------+
|Aggregated                                         |
+---------------------------------------------------+
|tcp: 8231, icmp: 7314, udp: 5523, igmp: 4423, egp: 2331|
+---------------------------------------------------+

总结一下,Spark提供的lower、struct、collect_list、map_from_entries等内置函数已经完全覆盖了你的需求,不需要额外写自定义UDF。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:16:03