【发布时间】:2019-10-03 00:51:03
【问题描述】:
我是使用 spark-sql 的 spark(2.x.x) 新手,我使用 spark sql 上下文创建了一个数据框。
dff = sqlCtx.read.format("com.databricks.spark.csv").option("header", "true").option("inferSchema","true").option("delimiter","\t").load("/home/hduser/Desktop/allFromDesktop/pyspark/creditData.csv")
dff.show()
Income|Limit|Rating|Cards|Age|Education|Gender|Student|Married| Ethnicity|Balance|Age_class|
+---+------------------+-----+------+-----+---+---------+------+-------+-------+----------------+-------+---------+
| 0|14.890999999999998| 3606| 283| 2| 34| 11| Male| No| Yes| Caucasian| 333| 25-34|
| 1| 106.025| 6645| 483| 3| 82| 15|Female| Yes| Yes| Asian| 903| 65+|
| 2|104.59299999999999| 7075| 514| 4| 71| 11| Male| No| No| Asian| 580| 65+|
我尝试对代码进行单元测试
tab=dff.select(['Age_class','Balance','Limit']).groupby('Age_class').agg
(F.count('Limit') ,F.mean('Limit').alias('Limit_avg'),F.min('Limit').alias('Limit_min'),F.max('Limit').alias('Limit_max')).withColumn('total',sum(col('Limit')).
over(Window)).withColumn('Percent',col('Limit')*100/col('total')).drop(c
ol('total'))
tab.show()
发现只要 agg(),count 函数执行,原始列就会被新的 .allias 列名替换 单元测试 #1(python 代码的第一部分成功执行)
tab=dff.select(['Age_class','Balance','Limit']).groupby('Age_class').agg(F.count('Limit_count')
,F.mean('Limit').alias('Limit_avg'),
F.min('Limit').alias('Limit_min'),
F.max('Limit').alias('Limit_max'))
tab.show()
有输出:
output:
---------+------------+------------------+---------+---------+
|Age_class|count(Limit)| Limit_avg|Limit_min|Limit_max|
+---------+------------+------------------+---------+---------+
| 45-54| 65| 4836.630769230769| 855| 11200|
| <25| 11|3932.6363636363635| 2120| 6375|
| 55-64| 68| 4530.0| 1311| 11966|
##Here you can see i lost my original 'Limit' column ##
具有“限制”列的原始数据框被删除(为什么??)被新列替换,即当代码运行 第二部分时我丢失了原始数据框列,无法找到原始数据框数据框中的列。 苏.. 语句的其余部分,即 *.withColumn('total',sum(col('Limit'))。 over(Window)).withColumn('Percent',col('Limit')100/col('total')).drop(col('total'))
显示错误:
Py4JJavaError: An error occurred while calling o2104.withColumn.
: org.apache.spark.sql.AnalysisException: cannot resolve '`Limit`' given input columns: [Age_class, Limit_max, Limit_min, Limit_avg, count(Limit)];;
现在,当我将 groupby 子句更改为 -->>.groupby('Age_class','Limit') insted of .groupby('Age_class') 时,我的代码会以所需的结果执行** 问题 1:为什么我需要在 groupby() 子句中添加“限制”列??当我已经在 SELECT 语句中提到它时 问题 2:执行后,即使 iam 使用“groupby”,“Age_class”列也没有转换为组(bins),请参阅下面的预期结果表,我期待类似的东西??
+---------+----------+------------------+----------+----------+-------+
|age_class|Limit|count(Limit)|Limit_avg|Limit_min|Limit_max|Percentage
+---------+----------+------------------+----------+----------+-------+
|
45-54 |120|3183.0666666666666|338|12612|12.0||
<25 |150| 2970.733333333333|276|15672|15.0||
55-64 |56| 3493.660714285714|385|15945|5.6||
35-44 |254| 3403.771653543307|250|15857|25.4||
25-34 |397| 3298.823677581864|343|18424|39.7||
65+ |23|3210.1739130434785|571|14896|2.3|
+---------+----------+------------------+----------+----------+-------+
tab=dff.select(['Age_class','Balance','Limit']).groupby('Age_class','Limit').agg(F.count('Limit')
,F.mean('Limit').alias('Limit_avg'),
F.min('Limit').alias('Limit_min'),
F.max('Limit').alias('Limit_max')).withColumn('total',sum(col('Limit')).over(Window)).withColumn('Percent',col('Limit')*100/col('total')).drop(col('total'))
tab.show()
实际输出('Age_class' 未转换成组(a.k.a bins)):
+---------+-----+------------+---------+---------+---------+-------------------+
|Age_class|Limit|count(Limit)|Limit_avg|Limit_min|Limit_max| Percent|
+---------+-----+------------+---------+---------+---------+-------------------+
| 45-54| 7838| 1| 7838.0| 7838| 7838| 0.4137807247233719|
| 35-44| 886| 1| 886.0| 886| 886|0.04677337612974069|
| 45-54| 4632| 1| 4632.0| 4632| 4632| 0.244530788073317|
| 55-64| 1448| 1| 1448.0| 1448| 1448|0.07644226708336853|
| 55-64| 5107| 1| 5107.0| 5107| 5107| 0.2696068080074331|
| 45-54| 2586| 1| 2586.0| 2586| 2586| 0.1365191316834192|
| 35-44| 4159| 1| 4159.0| 4159| 4159| 0.2195603513810288|
| 45-54| 4943| 1| 4943.0| 4943| 4943| 0.2609489821775488|
| 45-54| 2558| 1| 2558.0| 2558| 2558|0.13504096629782922|
| 25-34| 3969| 1| 3969.0| 3969| 3969|0.20952994340738237|
| 35-44| 5319| 1| 5319.0| 5319| 5319| 0.2807986316411859|
| 45-54| 8100| 1| 8100.0| 8100| 8100| 0.4276121294028212|
| 45-54| 6040| 1| 6040.0| 6040| 6040| 0.3188613903201284|
| 45-54| 4673| 1| 4673.0| 4673| 4673|0.24669524453078806|
| 65+| 2330| 1| 2330.0| 2330| 2330| 0.1230044767294535|
| 45-54| 6922| 1| 6922.0| 6922| 6922| 0.3654235999662134|
| 65+| 4263| 1| 4263.0| 4263| 4263|0.22505067995607736|
| 25-34| 4391| 1| 4391.0| 4391| 4391|0.23180800743306024|
| 65+| 7499| 1| 7499.0| 7499| 7499|0.39588436523355014|
| 45-54| 8732| 1| 8732.0| 8732| 8732| 0.460976433820424|
+---------+-----+------------+---------+---------+---------+-------------------+
【问题讨论】:
-
因为您按
Age_class分组。你想同时分组吗? -
不仅是 Age_class。
-
@vikassingh:请不要添加新问题(尤其是当有人已经回答时,答案会过时)。但是在这种情况下,它没有按您想要的方式分组的原因是您在两列上分组,
Age_class和Limit。您说您想在 groupnig 之后保留Limit行,但新的数据框(仅按Age_class分组)将有不同的行数。您需要考虑要保留Limit的哪个值。你已经有了最小值、最大值和平均值。
标签: apache-spark dataframe apache-spark-sql pyspark-sql