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

PySpark读取CSV时覆盖默认时间戳格式,为何不符合格式的记录未被拦截?

Spark CSV读取时不符合格式的时间戳未被标记为脏数据的原因

问题背景

测试CSV数据

ID|TIMESTAMP_COL
1|03-02-2003 08:37:55.671 PM
2|2003-02-03 08:37:55.671 AM

读取代码

from pyspark.sql.types import *

sch = StructType([StructField("ID",StringType(),False),StructField("TIMESTAMP_COL",StringType(),True)])

df = spark.read \
            .format("csv") \
            .option("encoding", "utf-8") \
            .option("mode", "PERMISSIVE") \
            .option("header", "true") \
            .option("dateFormat", "dd-MM-yyyy") \
            .option("timestampFormat", "dd-MM-yyyy HH:mm:ss.SSS a") \
            .option("delimiter", "|") \
            .option("columnNameOfCorruptRecord", "_corrupt_record") \
            .schema(sch) \
            .load("data.csv")

实际输出

df.show(truncate=False)

+-------------+-----------------------+-------------------+
|           ID|          TIMESTAMP_COL|    _corrupt_record|
+-------------+-----------------------+-------------------+
|            1|2003-02-03 08:37:55.671|               null|
|            2|0008-07-26 08:37:55.671|               null|
+-------------+-----------------------+-------------------+

按照设置的timestampFormat,ID为2的记录格式(yyyy-MM-dd)与指定格式(dd-MM-yyyy)不符,理应被标记为脏数据存入_corrupt_record,但实际却被解析成错误的时间值且未被标记,原因如下:

核心原因

  1. 字段类型不匹配
    你定义的Schema中,TIMESTAMP_COL被设置为StringType()而非TimestampType()。Spark的CSV读取器仅当字段类型为TimestampType()/DateType()时,才会严格按照timestampFormat/dateFormat校验格式;如果是字符串类型,Spark只会尝试按指定格式解析字符串为时间戳后再转回字符串,不会触发脏数据判定逻辑——因为字段本身是字符串,只要原始内容能被读取就不会被标记为脏数据。

  2. PERMISSIVE模式的触发条件
    PERMISSIVE模式仅在数据不符合Schema的类型要求(比如字符串转整数失败)时,才会将整条记录存入_corrupt_record。此处TIMESTAMP_COL是字符串类型,无论解析结果对错,原始字符串都能被正常读取,因此不会触发脏数据标记。而错误的解析结果是Spark对不匹配格式做容错解析导致的:它把2003-02-03中的2003当作日、02当作月、03当作年,经过日期计算后得到了错误的0008-07-26。

解决方法

  • 修改Schema类型:将TIMESTAMP_COL改为TimestampType(),这样Spark会严格校验格式,不符合的记录会被标记到_corrupt_record:
    sch = StructType([
        StructField("ID", StringType(), False),
        StructField("TIMESTAMP_COL", TimestampType(), True)
    ])
    
  • 读取后手动校验:保持原Schema,读取后用to_timestamp函数结合指定格式判断转换结果,过滤无效记录:
    from pyspark.sql.functions import to_timestamp
    
    df = df.withColumn(
        "valid_timestamp",
        to_timestamp("TIMESTAMP_COL", "dd-MM-yyyy HH:mm:ss.SSS a")
    ).filter("valid_timestamp IS NOT NULL")
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 20:45:27