使用PySpark API查询结果格式异常,Spark SQL无此问题
问题描述
执行Spark SQL查询航班延误数据时,返回结果格式正常:
spark.sql("""SELECT delay, origin, destination, CASE WHEN delay > 360 THEN 'Very Long Delays' WHEN delay > 120 AND delay < 360 THEN 'Long Delays' WHEN delay > 60 AND delay < 120 THEN 'Short Delays' WHEN delay > 0 and delay < 60 THEN 'Tolerable Delays' WHEN delay = 0 THEN 'No Delays' ELSE 'Early' END AS Flight_Delays FROM us_delay_flights_tbl ORDER BY origin, delay DESC""").show(10)
但执行等价逻辑的PySpark DataFrame API语句时,结果内容与SQL查询一致,但显示格式异常:
(df.select("delay","origin",col("destination"), when(df.delay > 360,"Very Long Delays") .when((df.delay > 120) & (df.delay < 360),"Long Delays") .when((df.delay > 60) & (df.delay < 120),"Short Delays") .when((df.delay > 0) & (df.delay < 60),"Tolerable Delays") .when((df.delay == 0),"No Delays") .otherwise("Early") ) .orderBy(asc("ORIGIN"),desc("delay")) ).show(10)
原因分析
- 未给动态生成的列命名:SQL语句通过
AS Flight_Delays给CASE表达式生成的列指定了明确的短名称;而DataFrame API中when链式调用生成的列没有设置别名,默认会显示为冗长的表达式字符串(如CASE WHEN ... ELSE ... END),过长的列名会破坏表格的显示对齐格式。 - 列名大小写不一致:
orderBy中使用了大写的"ORIGIN",但原DataFrame的列名是小写的"origin",虽然Spark默认对列名大小写不敏感,但这种不一致可能导致元数据处理时出现细微差异,间接影响显示格式。
解决方法
- 给新列添加别名:使用
alias()方法为when表达式生成的列指定与SQL一致的名称Flight_Delays,避免冗长的默认列名。 - 统一列名大小写:
orderBy中使用与原列一致的小写"origin",消除不必要的元数据处理差异。
修改后的代码如下:
(df.select("delay", "origin", col("destination"), when(df.delay > 360, "Very Long Delays") .when((df.delay > 120) & (df.delay < 360), "Long Delays") .when((df.delay > 60) & (df.delay < 120), "Short Delays") .when((df.delay > 0) & (df.delay < 60), "Tolerable Delays") .when(df.delay == 0, "No Delays") .otherwise("Early") .alias("Flight_Delays") ) .orderBy(asc("origin"), desc("delay")) ).show(10)
内容的提问来源于stack exchange,提问作者An old man in the sea.
相关产品推荐
相关产品推荐

