Spark合并两个DataFrame执行union时如何保留各自原始时间戳
Spark DataFrame Union操作后时间戳被统一覆盖问题修复
问题现象
对两个携带不同时间戳值的Spark DataFrame执行Union操作时,无法保留各自的原始时间戳,所有行最终返回统一的异常时间戳,不符合保留原始时间戳的预期:
- DataFrame1原始时间戳:
2022-07-8T05:08:22.395+000 - DataFrame2原始时间戳:
2022-07-8T05:02:10.757+000 - Union输出统一异常时间戳:
2022-07-8T05:08:34.651+000
当前使用Spark 3.1.2版本,两种Union实现方式均复现相同异常:
- 自定义补全缺失列的
unionMissing函数
def unionMissing(df1, df2=hist_df): # 为df1补充缺失列 left_df = df1 for column in set(df2.columns) - set(df1.columns): left_df = left_df.withColumn(column, F.lit(None)) # 为df2补充缺失列 right_df = df2 for column in set(df1.columns) - set(df2.columns): right_df = right_df.withColumn(column, F.lit(None)) # 对齐列顺序后执行union return left_df.unionAll(right_df.select(left_df.columns))
- 直接调用Spark内置
unionByNameAPI
df = df1.unionByName(df2, allowMissingColumns=True)
问题根因
该异常和Union算子本身逻辑无关,核心诱因是时间戳列使用非确定性函数生成,且未对中间计算结果做链路截断:
- 若时间戳列通过
F.current_timestamp()、F.now()这类非确定性内置函数或自定义UDF生成,Spark默认不会保存函数初次计算的结果,会在每次触发action作业时重新计算该列值。 - Union是典型的惰性转换算子,仅构建执行计划不会触发实际计算。当最终对Union后的结果执行
show()、write()等action时,Spark会回溯计算两个源DataFrame的时间戳列,叠加执行计划优化阶段的常量折叠、全阶段代码生成合并逻辑,最终所有行的时间戳会被统一赋值为action触发时刻计算出的同一个值。 - 输出结果中统一的时间戳
2022-07-8T05:08:34.651+000,本质就是触发Union结果计算时的集群当前时间,和两个源DataFrame的原始生成时间无关联。
修复方案
按优先级选择以下任意一种方案即可解决问题:
- Union前对源DataFrame做缓存截断执行链路
两个源DataFrame生成时间戳列后,立刻调用cache()/persist()并触发一次轻量action,截断非确定性函数的计算链路,强制保留生成时刻的时间戳值:# 生成df1、df2后立刻缓存并触发计算 df1 = df1.cache() df2 = df2.cache() df1.count() df2.count() # 缓存完成后再执行Union操作 df = df1.unionByName(df2, allowMissingColumns=True) - 生成时间戳时直接传入固定字面量
提前在Driver端获取固定时间值,包装为字面量列传入,避免Spark在后续Task计算中重复求值:from datetime import datetime # 错误写法:非确定性函数会被重复计算 # df1 = df1.withColumn("ts", F.current_timestamp()) # 正确写法:使用固定字面量赋值 fixed_ts = F.lit(datetime.utcnow()) df1 = df1.withColumn("ts", fixed_ts) df2 = df2.withColumn("ts", fixed_ts) - 临时关闭优化配置(仅调试场景使用,不推荐生产环境开启)
若临时调试不想修改代码,可在提交作业时添加以下配置规避优化逻辑:spark.sql.adaptive.enabled=false spark.sql.optimizer.constantFolding.enabled=false
内容的提问来源于stack exchange,提问作者Jeevan Kande
相关产品推荐
相关产品推荐

