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
相关产品推荐
相关产品推荐

