Databricks Autoloader写入Delta表失败,报错long overflow求助
问题:Autoloader写入Delta表触发Long溢出错误
代码片段
df.writeStream.format("delta").foreachBatch(lambda df, epochId: update_insert(df, epochId, cdm)).option("checkpointLocation", checkpoint_directory).trigger(availableNow=True).start()
报错堆栈
org.apache.spark.SparkException: Job aborted due to stage failure: Task 4 in stage 649.0 failed 4 times, most recent failure: Lost task 4.3 in stage 649.0 (TID 2611) (172.20.12.5 executor 3): org.apache.spark.SparkException: Task failed while writing rows. at com.databricks.photon.PhotonWriteStageExec.$anonfun$executeWrite$2(PhotonWriteStageExec.scala:132) at com.databricks.photon.PhotonExec.$anonfun$executePhoton$6(PhotonExec.scala:579) at com.databricks.photon.PhotonExec.$anonfun$executePhoton$6$adapted(PhotonExec.scala:422) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsWithIndexInternal$2(RDD.scala:916) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsWithIndexInternal$2$adapted(RDD.scala:916) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:60) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:406) at org.apache.spark.rdd.RDD.iterator(RDD.scala:370) at org.apache.spark.scheduler.ResultTask.$anonfun$runTask$3(ResultTask.scala:75) at com.databricks.spark.util.ExecutorFrameProfiler$.record(ExecutorFrameProfiler.scala:110) at org.apache.spark.scheduler.ResultTask.$anonfun$runTask$1(ResultTask.scala:75) at com.databricks.spark.util.ExecutorFrameProfiler$.record(ExecutorFrameProfiler.scala:110) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:55) at org.apache.spark.scheduler.Task.doRunTask(Task.scala:179) at org.apache.spark.scheduler.Task.$anonfun$run$5(Task.scala:142) at com.databricks.unity.UCSEphemeralState$Handle.runWith(UCSEphemeralState.scala:41) at com.databricks.unity.HandleImpl.runWith(UCSHandle.scala:99) at com.databricks.unity.HandleImpl.$anonfun$runWithAndClose$1(UCSHandle.scala:104) at scala.util.Using$.resource(Using.scala:269) at com.databricks.unity.HandleImpl.runWithAndClose(UCSHandle.scala:103) at org.apache.spark.scheduler.Task.$anonfun$run$1(Task.scala:142) at com.databricks.spark.util.ExecutorFrameProfiler$.record(ExecutorFrameProfiler.scala:110) at org.apache.spark.scheduler.Task.run(Task.scala:97) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$13(Executor.scala:904) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1713) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:907) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at com.databricks.spark.util.ExecutorFrameProfiler$.record(ExecutorFrameProfiler.scala:110) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:761) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) Caused by: java.lang.ArithmeticException: long overflow at java.lang.Math.multiplyExact(Math.java:892) at org.apache.spark.sql.catalyst.util.DateTimeUtils$.millisToMicros(DateTimeUtils.scala:257) at org.apache.spark.sql.catalyst.util.RebaseDateTime$.rebaseGregorianToJulianMicros(RebaseDateTime.scala:374) at org.apache.spark.sql.catalyst.util.RebaseDateTime$.rebaseGregorianToJulianMicros(RebaseDateTime.scala:394) at org.apache.spark.sql.catalyst.util.RebaseDateTime.rebaseGregorianToJulianMicros(RebaseDateTime.scala)
问题根源
溢出发生在Spark的日期时间转换过程中:DateTimeUtils$.millisToMicros方法将毫秒数乘以1000转成微秒时,超出了Java long类型的最大值(9223372036854775807)。这是因为源数据中存在非法日期时间值(比如远早于1900年、远超未来的时间戳,或格式错误被解析成异常数值)。
解决方案
- 过滤异常日期值
在流处理的前置步骤中添加校验,过滤掉超出合法范围的日期记录:
from pyspark.sql.functions import col, to_timestamp # 示例:限制日期在1900-2100年之间 valid_df = df.filter(col("your_timestamp_column").between(to_timestamp("1900-01-01"), to_timestamp("2100-12-31"))) # 或直接校验毫秒数范围(对应1900-01-01到2100-12-31) valid_df = df.filter((col("your_timestamp_column").cast("long") >= -2208988800000) & (col("your_timestamp_column").cast("long") <= 4102444799000))
- 调整Spark日期重基配置
禁用可能触发溢出的日期转换逻辑,在Spark会话初始化时添加以下配置:
spark.conf.set("spark.sql.datetimeRebaseModeInWrite", "CORRECTED") spark.conf.set("spark.sql.parquet.datetimeRebaseModeInWrite", "CORRECTED")
- 检查自定义函数逻辑
排查update_insert函数中是否有手动处理日期的代码,比如手动计算毫秒/微秒的操作,替换为Spark内置日期函数(如to_timestamp、date_format)避免手动计算溢出。
内容的提问来源于stack exchange,提问作者Greencolor
相关产品推荐
相关产品推荐

