PySpark DataFrame如何按时间段筛选行及查询异常原因排查
问题解答
时间提取格式验证
你的Pickup_time提取逻辑是正确的:date_format('Pickup_datetime', 'HH:mm:ss')可以从Timestamp类型的原始字段中,正确提取出24小时制的时分秒字符串,从你上传的表结构截图也能验证该列格式符合预期。
查询报错的原因
你写的查询代码存在3处核心错误:
- 时间字符串未加引号:筛选条件中的
01:00:00和06:00:00不是合法的Python/Spark语法元素,需要用引号包裹为字符串格式'01:00:00'、'06:00:00' - 筛选与分组顺序颠倒:
where过滤逻辑需要放在groupBy之前执行,你当前的写法是先分组再过滤,既不符合业务逻辑(应该先筛出夜间订单再统计),也会触发Spark语法报错 - 列引用风险:链式调用中直接用
df.Pickup_time引用列可能出现指向错误,建议统一用f.col("列名")的方式引用字段
修正后的可用代码
你不需要额外创建COUNT常量列,直接用count聚合即可,优化后的代码如下:
import pyspark.sql.functions as f result_df = df.where( (f.col("Pickup_time") >= '01:00:00') & (f.col("Pickup_time") <= '06:00:00') ).groupBy("Pickup_date", "Driver_ID") \ .agg(f.count("*").alias("Total_Rides")) \ .orderBy("Pickup_date", ascending=False)
内容的提问来源于stack exchange,提问作者Adetya Aggarwal
相关产品推荐
相关产品推荐

