如何用PySpark对吉他厂商的吉他类型数据分组统计并排序?
用PySpark实现吉他厂商的吉他类型统计与排序
需求说明
输入CSV数据包含vendor(吉他厂商)和product(吉他产品)两个字段,产品格式为吉他类型::具体型号,需要完成以下操作:
- 按厂商统计每种吉他类型的数量
- 按
acoustic类型的数量降序排序 - 输出格式为
厂商 - acoustic:数量, electric:数量
实现代码与步骤
1. 初始化SparkSession并读取CSV数据
from pyspark.sql import SparkSession from pyspark.sql.functions import split, count, col, lit, concat_ws, desc # 初始化SparkSession spark = SparkSession.builder.appName("GuitarVendorStats").getOrCreate() # 读取CSV文件,自动识别表头和数据类型 df = spark.read.csv("guitars.csv", header=True, inferSchema=True)
2. 提取吉他类型字段
通过拆分product字段,分离出吉他类型(acoustic或electric):
df_with_type = df.withColumn("guitar_type", split(col("product"), "::").getItem(0))
3. 分组统计各类型数量
先按厂商和吉他类型分组统计数量,再通过透视表将吉他类型转为列,方便后续排序和格式化:
# 分组统计每个厂商下各吉他类型的数量 grouped_df = df_with_type.groupBy("vendor", "guitar_type").agg(count("*").alias("count")) # 透视表转换:将吉他类型作为列,对应数量为值 pivoted_df = grouped_df.groupBy("vendor").pivot("guitar_type").sum("count") # 对缺失的类型补0(比如某厂商没有生产某类型吉他时) pivoted_df = pivoted_df.fillna(0, subset=["acoustic", "electric"])
4. 按acoustic数量降序排序
以acoustic列的数值作为排序依据,降序排列:
sorted_df = pivoted_df.orderBy(desc("acoustic"))
5. 格式化并输出结果
将数据拼接成要求的字符串格式,然后收集结果打印:
# 拼接成目标输出格式 result_df = sorted_df.withColumn( "result", concat_ws( " - ", col("vendor"), concat_ws(", ", concat_ws(":", lit("acoustic"), col("acoustic")), concat_ws(":", lit("electric"), col("electric")) ) ) ) # 打印结果 for row in result_df.select("result").collect(): print(row.result)
运行上述代码后,将得到如下输出:
fender - acoustic:2, electric:1 yamaha - acoustic:1, electric:1
内容的提问来源于stack exchange,提问作者GuyMatz
相关产品推荐
相关产品推荐

