Spark DataFrame转换需求:按分组统计值出现频次并转宽表
Spark DataFrame长表转宽表(统计值出现次数)
需求说明
将包含分组列和值列的长表,转换为以分组列为行、不同值为列,值为对应出现次数的宽表,无对应值的位置填充0。
原DataFrame结构
+----------+---------+ |Column 1 | Values | +----------+---------+ | A | value1 | | B | value2 | | C | value2 | | A | value1 | | B | value3 | | C | value1 | | A | value1 | | B | value1 | | C | value2 | +----------+---------+
实现方案
可以通过groupBy+pivot的方式快速实现,以下提供PySpark和Scala两种版本的代码:
PySpark 实现
from pyspark.sql import SparkSession from pyspark.sql.functions import count # 初始化SparkSession spark = SparkSession.builder.appName("pivot_count").getOrCreate() # 构造原DataFrame data = [ ("A", "value1"), ("B", "value2"), ("C", "value2"), ("A", "value1"), ("B", "value3"), ("C", "value1"), ("A", "value1"), ("B", "value1"), ("C", "value2") ] df = spark.createDataFrame(data, ["Column 1", "Values"]) # 直接分组、透视并统计次数,填充空值为0 result_df = df.groupBy("Column 1").pivot("Values").agg(count("*")).fillna(0) # 查看结果 result_df.show()
Scala 实现
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.count object PivotCountExample { def main(args: Array[String]): Unit = { val spark = SparkSession.builder.appName("pivot_count").getOrCreate() import spark.implicits._ // 构造原DataFrame val data = Seq( ("A", "value1"), ("B", "value2"), ("C", "value2"), ("A", "value1"), ("B", "value3"), ("C", "value1"), ("A", "value1"), ("B", "value1"), ("C", "value2") ) val df = data.toDF("Column 1", "Values") // 分组、透视统计次数,填充空值为0 val resultDF = df.groupBy("Column 1").pivot("Values").agg(count("*")).na.fill(0) // 查看结果 resultDF.show() } }
代码说明
- groupBy("Column 1"):按分组列聚合数据
- pivot("Values"):将
Values列的不同取值转换为新的列 - agg(count("*")):统计每个分组下各值的出现次数
- fillna(0)/na.fill(0):将没有对应值的单元格填充为0
执行后得到的结果如下:
+----------+-------+-------+-------+ |Column 1 |value1 |value2 |value3 | +----------+-------+-------+-------+ | A | 3| 0| 0| | B | 1| 1| 1| | C | 1| 2| 0| +----------+-------+-------+-------+
内容的提问来源于stack exchange,提问作者Error-F
相关产品推荐
相关产品推荐

