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

PySpark处理早于1970年日期触发OverflowError的解决方法

问题:处理早于1970-01-01的日期时mapPartitions转DataFrame报错OverflowError

在Mac OS Monterey v12.6.1本地环境使用PySpark 3.3.0处理包含早于1970-01-01的日期和时间戳字段的Parquet文件时,读取文件、写入Postgres表以及直接写入Parquet文件均正常,但通过mapPartitions捕获Postgres插入失败的记录并转为DataFrame写入时,触发以下错误:

22/11/27 20:22:46 ERROR Utils: Aborting task
org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/worker.py", line 686, in main
    process()
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/worker.py", line 678, in process
    serializer.dump_stream(out_iter, outfile)
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/serializers.py", line 273, in dump_stream
    vs = list(itertools.islice(iterator, batch))
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/util.py", line 81, in wrapper
    return f(*args, **kwargs)
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/sql/types.py", line 788, in toInternal
    return tuple(
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/sql/types.py", line 789, in <genexpr>
    f.toInternal(v) if c else v
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/sql/types.py", line 591, in toInternal
    return self.dataType.toInternal(obj)
  File "/Users/pm/opt/spark-3.3.0-bin-hadoop3/python/lib/pyspark.zip/pyspark/sql/types.py", line 216, in toInternal
    calendar.timegm(dt.utctimetuple()) if dt.tzinfo else time.mktime(dt.timetuple())
OverflowError: mktime argument out of range

相关代码逻辑

mapPartitions核心逻辑

for row in rows:
    try:
        processed_row_count.add(1) # accumulator
        cur.execute(insert_statement, row)
        accepted_row_count.add(1) # accumulator
    except Exception as e:
        # Collect the records with error message here and write to reject file
        rejected_row_count.add(1) # accumulator
        row = row.asDict()
        row["ErrorMessage"] = f"Error received from psycopg2 module is: {str(e)}"
        yield Row(**row)
conn.close()

调用mapPartitions的代码

rejected_on_load_df = enriched_df.rdd.mapPartitions(process_rows).toDF(
        enriched_df.schema.add("ErrorMessage", StringType(), nullable=False)
    )

解决方案

问题原因

错误出现在Spark将Python的datetime对象序列化为内部格式时,time.mktime()无法处理早于1970-01-01的本地时间(Unix时间戳以1970年为起点,部分系统的mktime实现不支持负数时间戳)。而直接读写Parquet/Postgres时,Spark使用内部日期时间处理逻辑,绕过了Python层的mktime调用,因此没有报错。

可行解决办法

方案1:将日期时间字段转为字符串处理

在生成失败记录的Row时,把早于1970-01-01的日期/时间戳字段转为字符串,避免触发mktime转换:

from datetime import datetime

for row in rows:
    try:
        processed_row_count.add(1)
        cur.execute(insert_statement, row)
        accepted_row_count.add(1)
    except Exception as e:
        rejected_row_count.add(1)
        row_dict = row.asDict()
        # 遍历字段,将早于1970年的datetime转为字符串
        epoch_start = datetime(1970, 1, 1)
        for key, value in row_dict.items():
            if isinstance(value, datetime) and value < epoch_start:
                row_dict[key] = value.strftime("%Y-%m-%d %H:%M:%S")
        row_dict["ErrorMessage"] = f"Error received from psycopg2 module is: {str(e)}"
        yield Row(**row_dict)
conn.close()

同时修改生成DataFrame的schema,将对应日期时间字段改为StringType():

from pyspark.sql.types import StringType

# 复制原schema并修改日期时间字段类型
new_schema = enriched_df.schema.copy()
# 根据实际字段名修改,示例假设字段名为"old_date"、"old_timestamp"
for field in new_schema.fields:
    if field.dataType.typeName() in ["date", "timestamp"]:
        new_schema = new_schema.withField(field.name, StringType())
# 添加ErrorMessage字段
new_schema = new_schema.add("ErrorMessage", StringType(), nullable=False)

rejected_on_load_df = enriched_df.rdd.mapPartitions(process_rows).toDF(new_schema)

方案2:使用Spark SQL方式捕获失败记录(推荐)

避开RDD的Python序列化问题,改用Spark原生write API结合批量处理逻辑:

from pyspark.sql.functions import lit

def write_to_postgres_and_capture_failures(df):
    try:
        # 尝试批量写入Postgres
        df.write \
          .format("jdbc") \
          .option("url", "jdbc:postgresql://host:port/db") \
          .option("dbtable", "target_table") \
          .option("user", "user") \
          .option("password", "password") \
          .option("batchsize", 1000) \
          .mode("append") \
          .save()
    except Exception as e:
        # 写入失败时,标记错误后写入拒绝存储
        df.withColumn("ErrorMessage", lit(f"Batch insert failed: {str(e)}")) \
          .write \
          .mode("append") \
          .parquet("/path/to/rejected_records")

# 批处理场景直接调用
write_to_postgres_and_capture_failures(enriched_df)

该方案利用Spark原生JDBC写入逻辑,避免了Python层的日期序列化问题,同时简化了错误捕获流程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 02:15:35