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

如何在PySpark/Pandas中校验季度数据完整性?

季度数据校验方案(Pandas + PySpark实现)

问题背景

现有结构如下的数据集:

process_date    ItemNo         ItemType
01-Mar-2019       1               abc
01-Jun-2019       2               cde
01-Sep-2019       1               abc

数据按季度交付,process_date为季度首日。需完成两项校验:

  1. 确认文件包含当前交付周期对应的目标季度数据(比如12月交付的文件必须包含当年Q3数据,仅到Q2则失败);
  2. 单独校验文件是否包含目标季度的上一季度数据(比如9月交付的文件需包含Q2数据)。

Pandas 实现

1. 日期预处理与季度提取

先把字符串格式的日期转成datetime类型,并生成年季度标识:

import pandas as pd

# 读取数据(替换为实际路径)
df = pd.read_csv("quarterly_data.csv")
df['process_date'] = pd.to_datetime(df['process_date'], format='%d-%b-%Y')
# 生成如"2019Q3"的季度标识
df['year_quarter'] = df['process_date'].dt.to_period('Q')

2. 校验目标季度是否存在

假设当前交付的文件对应目标季度为target_q(比如12月交付对应2019Q3),直接校验该季度是否在数据集中:

# 示例:设置目标季度
target_q = pd.Period('2019Q3', freq='Q')

# 校验逻辑
if target_q not in df['year_quarter'].unique():
    raise ValueError("文件未包含对应季度数据,处理终止")

3. 校验上一季度是否存在

基于目标季度计算上一季度,再检查存在性:

# 计算上一季度
prev_q = target_q - 1

# 校验逻辑
if prev_q not in df['year_quarter'].unique():
    raise ValueError("文件未包含上一季度数据")

PySpark 实现

1. 日期预处理与季度提取

将字符串日期转为Spark日期类型,并生成年季度字符串:

from pyspark.sql import SparkSession
from pyspark.sql.functions import date_format, expr

spark = SparkSession.builder.appName("QuarterCheck").getOrCreate()

# 读取数据(替换为实际路径)
df = spark.read.csv("quarterly_data.csv", header=True)
# 转换日期格式
df = df.withColumn("process_date", expr("to_date(process_date, 'dd-MMM-yyyy')"))
# 生成如"2019Q3"的季度标识
df = df.withColumn("year_quarter", date_format("process_date", "yyyy'Q'q"))

2. 校验目标季度是否存在

设置目标季度字符串,通过计数判断是否存在:

# 示例:设置目标季度
target_q = "2019Q3"

# 校验逻辑
if df.filter(df.year_quarter == target_q).count() == 0:
    raise ValueError("文件未包含对应季度数据,处理终止")

3. 校验上一季度是否存在

先写一个通用函数计算上一季度,再执行校验:

def get_previous_quarter(target_q):
    year = int(target_q[:4])
    quarter = int(target_q[-1])
    if quarter == 1:
        return f"{year-1}Q4"
    else:
        return f"{year}Q{quarter-1}"

# 计算上一季度并校验
prev_q = get_previous_quarter(target_q)
if df.filter(df.year_quarter == prev_q).count() == 0:
    raise ValueError("文件未包含上一季度数据")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 03:40:33