使用Databricks PySpark分析Uber数据时遇异常错误求助
问题排查:PySpark找出每个调度基地订单量最多的星期几
需求与问题
需求为找出每个调度基地(dispatching_base_number)订单量最多的星期几,执行以下代码时finaldf.show()报错:
from pyspark.sql import SparkSession from pyspark.sql.functions import expr, sum, rank from pyspark.sql.window import Window spark = SparkSession.builder \ .appName("first") \ .getOrCreate() df1 = spark.read.format("csv") \ .option("header", "true") \ .load("dbfs:/FileStore/tables/uber.csv") df2 = df1.withColumn("trips",col("trips").cast("integer")) finaldf = df2.withColumn("Day", expr("DATE_FORMAT(TO_DATE(DATE,'MM/dd/yyyy'),'EEE')")) \ .groupBy("dispatching_base_number", "Day") \ .agg(sum("trips")).alias("Sum") \ .withColumn("rnk", rank().over(Window.partitionBy("dispatching_base_number").orderBy("Sum"))) \ .filter("rnk = 1") \ .drop("rnk") finaldf.show()
错误原因与修正方案
1. 缺失函数导入
代码中使用了col()函数,但未从pyspark.sql.functions导入该函数,会触发NameError。需补充导入col及排序需要的desc函数。
2. 聚合列别名位置错误
原代码中.agg(sum("trips")).alias("Sum")写法错误,别名需直接绑定在聚合函数上,否则后续无法引用Sum列,正确写法为.agg(sum("trips").alias("Sum"))。
3. 排序方向错误
要筛选订单量最多的记录,窗口函数排序需使用降序,原代码升序排序会取到订单量最少的结果,需改为orderBy(desc("Sum"))。
修正后的代码
from pyspark.sql import SparkSession from pyspark.sql.functions import expr, sum, rank, col, desc from pyspark.sql.window import Window spark = SparkSession.builder \ .appName("first") \ .getOrCreate() # 读取CSV数据 df1 = spark.read.format("csv") \ .option("header", "true") \ .load("dbfs:/FileStore/tables/uber.csv") # 将trips列转为整数类型 df2 = df1.withColumn("trips", col("trips").cast("integer")) # 计算每个基地每天的订单总量,筛选每个基地订单最多的星期几 finaldf = df2.withColumn("Day", expr("DATE_FORMAT(TO_DATE(DATE,'MM/dd/yyyy'),'EEE')")) \ .groupBy("dispatching_base_number", "Day") \ .agg(sum("trips").alias("Sum")) \ .withColumn("rnk", rank().over(Window.partitionBy("dispatching_base_number").orderBy(desc("Sum")))) \ .filter("rnk = 1") \ .drop("rnk") finaldf.show()
额外说明
如果数据源中DATE列的格式与MM/dd/yyyy不匹配,需调整TO_DATE函数的格式参数,否则会生成null值影响计算结果。
内容的提问来源于stack exchange,提问作者suman
相关产品推荐
相关产品推荐

