如何使用PySpark的GroupBy找出总信用限额最少的收入类别
PySpark 解决方案:按收入类别统计总信用限额并找出最小值
你的代码问题分析
- 直接用
map提取单个字段得到的RDD没有键值对结构,join操作要求RDD是(键, 值)格式,所以你的join写法完全错误。 - 没有将
Credit_Limit从字符串转为数值类型,直接做求和会引发类型错误。 - 没必要拆分两个RDD再关联,直接从原数据中同时提取需要的字段即可。
推荐方案:使用DataFrame API(更简洁高效)
DataFrame API是PySpark的首选方式,语法直观且性能更优:
# 读取数据,自动推断字段类型(包括将Credit_Limit转为数值型) WHData = spark.read.option("header", True) \ .option("inferSchema", True) \ .csv("file:///home/prac/test3/input/CreditCard.csv") # 分组求和→重命名字段→排序取最小值 result = WHData.groupBy("Income_Category") \ .sum("Credit_Limit") \ .withColumnRenamed("sum(Credit_Limit)", "Total_Credit_Limit") \ .orderBy("Total_Credit_Limit") \ .first() print(f"总信用限额最少的收入类别:{result['Income_Category']},总计:{result['Total_Credit_Limit']}")
代码说明
inferSchema=True让Spark自动识别字段类型,省去手动转换Credit_Limit的步骤。groupBy+sum完成分组求和逻辑,withColumnRenamed优化字段可读性。orderBy+first直接获取排序后的第一条数据,即总信用限额最小的记录。
备选方案:使用RDD API(若考核要求必须用RDD)
如果需要用RDD实现,调整后的代码如下:
WHData = spark.read.option("header", True).csv("file:///home/prac/test3/input/CreditCard.csv") # 映射为(收入类别, 浮点型信用限额)的键值对RDD income_credit_rdd = WHData.rdd.map(lambda row: (row[4], float(row[7]))) # 分组求和→转本地列表→排序取最小值 total_credit_by_income = income_credit_rdd.reduceByKey(lambda a, b: a + b).collect() min_total = sorted(total_credit_by_income, key=lambda x: x[1])[0] print(f"总信用限额最少的收入类别:{min_total[0]},总计:{min_total[1]}")
代码说明
map将每条记录转换为符合RDD聚合要求的键值对结构,同时完成类型转换。reduceByKey按收入类别聚合,对信用限额求和。collect将分布式RDD结果转为本地列表,再通过sorted排序提取最小值。
内容的提问来源于stack exchange,提问作者Jay
相关产品推荐
相关产品推荐

