Delta Live Tables与PySpark聚合结果不一致问题求助
问题原因分析
核心问题是两段代码的分组字段完全不一致:
- PySpark查询中使用
groupBy("job_title"),按职位标题维度进行分组聚合 - DLT流水线中使用
groupBy("job_industry_category"),按职位所属行业分类维度进行分组聚合
分组维度不同,必然会得到差异极大的聚合结果,数据模式(列结构)也会因为分组键的不同而出现明显区别——PySpark结果的行对应各个职位标题,DLT结果的行对应各个行业分类,即便pivot("owns_car")生成的列(如Yes/No)可能一致,但每行的统计值对应的维度完全不同。
另外可以额外确认两个查询的源数据一致性:检查PySpark中的filter_df和DLT中的filtered_customers是否是过滤逻辑完全相同的数据集,如果源数据本身存在差异,也会影响结果,但从代码来看,分组字段写错是最直接的原因。
将DLT代码中的groupBy("job_industry_category")修改为groupBy("job_title"),重新运行流水线后,结果应该会和PySpark查询一致。
内容的提问来源于stack exchange,提问作者awesome_sangram
相关产品推荐
相关产品推荐

