You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark查询结果与同数据集SQL查询结果不符问题排查

PySpark与SQL查询结果不一致的排查方案

可能的原因

  • 日期格式符错误:PySpark查询中使用了小写的mm作为月份格式符,但Spark的日期格式规则里,mm代表分钟,正确的月份格式符是大写MM。这会导致to_date函数解析CallDate时出错,过滤条件year(...) == 2018无法正确筛选出2018年的数据,结果混入了其他年份的高Delay记录;而SQL查询使用了正确的yyyy-MM-dd格式,过滤逻辑正常。
  • 数据源不一致:PySpark使用的fire_df与SQL查询的demo_db.fire_service_calls_tbl可能经过不同预处理(比如fire_df未清洗异常值,而表中已过滤),导致Delay列的数据范围不同。
  • 数据类型差异:Delay列在fire_df和数据库表中的数据类型不一致(比如一个是字符串、一个是数值),会导致排序逻辑出现偏差。

调试步骤

  1. 修正日期格式符:将PySpark查询中的yyyy-mm-dd改为yyyy-MM-dd,重新运行后对比结果:

    san_francisco_worst_response_time = fire_df.where(year(to_date(col('CallDate'), 'yyyy-MM-dd')) == 2018) \
                                                .select('Neighborhood', 'Delay') \
                                                .orderBy('Delay', ascending=False)
    display(san_francisco_worst_response_time)
    
  2. 验证过滤结果一致性:分别统计两者过滤后的记录数和Delay极值:

    • PySpark端:
      # 统计记录数
      print(f"PySpark过滤后记录数: {fire_df.where(year(to_date(col('CallDate'), 'yyyy-MM-dd')) == 2018).count()}")
      # 查看Delay的最大/最小值
      fire_df.where(year(to_date(col('CallDate'), 'yyyy-MM-dd')) == 2018).select(max('Delay'), min('Delay')).show()
      
    • SQL端:
      SELECT COUNT(*), MAX(Delay), MIN(Delay) 
      FROM demo_db.fire_service_calls_tbl
      WHERE YEAR(TO_DATE(CallDate, "yyyy-MM-dd")) = 2018;
      

    对比两组结果,确认记录数和极值是否匹配。

  3. 核对数据源一致性:确认fire_df是直接读取原始CSV文件未做修改,而demo_db.fire_service_calls_tbl是否基于同一文件创建,中间没有额外的清洗逻辑(比如过滤Delay大于100的记录)。

  4. 检查列数据类型:查看Delay列的类型是否一致:

    • PySpark端:
      fire_df.printSchema()
      
    • SQL端:
      DESCRIBE demo_db.fire_service_calls_tbl;
      

    确保两者的Delay列均为数值类型(如double),避免因类型差异导致排序异常。

  5. 排查异常记录:如果修正日期格式后仍不一致,提取PySpark结果中高Delay值的CallDate,检查这些记录是否存在于SQL表中,是否被表的预处理逻辑排除。

内容的提问来源于stack exchange,提问作者Rahul Saran

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 12:13:08