PySpark季度数据缺失检测:整合跨年逻辑实现一行代码
PySpark季度数据缺失检测逻辑整合方案
问题场景
我有一个存储季度数据的PySpark DataFrame,数据示例如下:
2022-03-01 abc 2022-06-01 xyz 2000-03-01 abcd
校验需求
- 检测1960年至今的季度数据缺失情况
- 当前年份仅校验已过季度(例如2022年仅检查前3个季度)
- 1965年为例外,无需校验其全年数据
现有代码仅能处理非当前年份的校验,需要将当前年份的校验逻辑整合到现有代码中,实现一行链式调用完成所有校验。现有代码如下:
qtrs = df.groupBy(year("mydate").alias("q_count")).count().filter(col("count")!= 4).filter(~col("qtr_count").isin(1965)).collect() if len(qtrs) !=0: return ("Error")
整合后的解决方案
可以实现逻辑整合,核心是针对当前年份动态计算应有的已过季度数,结合往年校验逻辑通过链式调用完成。修改后的代码如下:
from pyspark.sql.functions import year, quarter, current_date, ceil, month, countDistinct, col qtrs = df.groupBy(year("mydate").alias("year")) \ .agg(countDistinct(quarter("mydate")).alias("actual_qtrs")) \ .filter( (~col("year").isin(1965)) & ( (col("year") < year(current_date())) & (col("actual_qtrs") != 4) | (col("year") == year(current_date())) & (col("actual_qtrs") != ceil(month(current_date())/3)) ) ).collect() if len(qtrs) != 0: return "Error"
关键逻辑说明
- 使用
countDistinct(quarter("mydate"))统计每年实际存在的唯一季度数,避免重复数据导致的统计误差 - 过滤条件分两类场景:
- 非当前年份且非1965年:校验实际季度数是否不等于4(全年应有的季度数)
- 当前年份:通过
ceil(month(current_date())/3)计算已过季度数(例如9月对应3,12月对应4),校验实际季度数是否与该值不符
- 所有校验逻辑通过链式调用整合到一个
filter操作中,满足"一行代码"的需求
内容的提问来源于stack exchange,提问作者Pysparker
相关产品推荐
相关产品推荐

