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和数据库表中的数据类型不一致(比如一个是字符串、一个是数值),会导致排序逻辑出现偏差。
调试步骤
修正日期格式符:将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)验证过滤结果一致性:分别统计两者过滤后的记录数和
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;
对比两组结果,确认记录数和极值是否匹配。
- PySpark端:
核对数据源一致性:确认
fire_df是直接读取原始CSV文件未做修改,而demo_db.fire_service_calls_tbl是否基于同一文件创建,中间没有额外的清洗逻辑(比如过滤Delay大于100的记录)。检查列数据类型:查看
Delay列的类型是否一致:- PySpark端:
fire_df.printSchema() - SQL端:
DESCRIBE demo_db.fire_service_calls_tbl;
确保两者的
Delay列均为数值类型(如double),避免因类型差异导致排序异常。- PySpark端:
排查异常记录:如果修正日期格式后仍不一致,提取PySpark结果中高Delay值的
CallDate,检查这些记录是否存在于SQL表中,是否被表的预处理逻辑排除。
内容的提问来源于stack exchange,提问作者Rahul Saran
相关产品推荐
相关产品推荐

