PySpark调用write经JDBC写入Azure时DataFrame数据异常变更问题
问题描述
使用PySpark通过JDBC将DataFrame写入Azure数据库时出现异常:执行write操作过程中DataFrame数据无理由变更,向Azure端写入了错误数据,写入操作完成后原PySpark DataFrame也留存了同一份错误数据。
核心代码如下:
print(sparkDF_cleaned.show()) sparkDF_cleaned.write \ .format("jdbc") \ .mode("overwrite") .option("url", jdbcUrl) \ .option("dbtable", "dbo.upsert_test") \ .option("user", jdbcUsername) \ .option("password", jdbcPassword) \ .save() print(f"data loaded to table {db_table_name}") print(sparkDF_cleaned.show())
代码执行输出如下:
sparkDF_cleaned : +------------+---+-----+----------+----------------------+ | id_date| id|value| _date|datetime_of_extraction| +------------+---+-----+----------+----------------------+ |1 2022-05-01| 1| 17|2022-05-01| 2022-06-01| |1 2022-05-06| 1| 6|2022-05-06| 2022-06-13| |2 2022-05-02| 2| 10|2022-05-02| 2022-06-01| |3 2022-05-03| 3| 15|2022-05-03| 2022-06-01| +------------+---+-----+----------+----------------------+ data loaded to table upsert_test sparkDF_cleaned : +------------+---+-----+----------+----------------------+ | id_date| id|value| _date|datetime_of_extraction| +------------+---+-----+----------+----------------------+ |1 2022-05-06| 1| 6|2022-05-06| 2022-06-13| |2 2022-05-02| 2| 5|2022-05-02| 2022-06-13| +------------+---+-----+----------+----------------------+
最终Azure表中接收到的就是第二次打印的错误数据。
问题原因
核心原因有两点:
- PySpark DataFrame 惰性求值机制,未做结果缓存:
show()、write.save()都属于Action算子,每次调用时Spark都会沿着DataFrame的血缘链路从头重新计算结果,不会复用之前Action触发计算得到的数据。如果sparkDF_cleaned的上游计算逻辑包含非确定性操作——比如使用了current_timestamp()这类实时取系统时间的函数、随机数函数,或是上游依赖的数据源在任务执行期间被其他操作修改,就会出现两次Action得到的数据完全不一致的情况,write操作写入的就是重算后的错误数据。 - 代码存在语法瑕疵:
.mode("overwrite")行末尾缺失了\行续接符,导致write链路的链式调用断裂,部分场景下会触发执行计划解析异常,放大数据重算的不一致问题。
修复方案
- 对需要固定结果的DataFrame提前做缓存,在第一次触发Action前调用
cache()/persist()方法,将计算结果固化到存储中,后续所有Action都会复用缓存数据,不会触发上游重算。示例代码:
# 缓存DataFrame,固化计算结果 sparkDF_cleaned = sparkDF_cleaned.cache() # 第一次查看数据 print(sparkDF_cleaned.show()) # 写入操作,补全缺失的行续接符 sparkDF_cleaned.write \ .format("jdbc") \ .mode("overwrite") \ .option("url", jdbcUrl) \ .option("dbtable", "dbo.upsert_test") \ .option("user", jdbcUsername) \ .option("password", jdbcPassword) \ .save() print(f"data loaded to table dbo.upsert_test") # 第二次查看数据,结果和第一次完全一致 print(sparkDF_cleaned.show()) # 用完后释放缓存资源 sparkDF_cleaned.unpersist()
- 排查
sparkDF_cleaned的上游计算逻辑,替换非确定性操作:如果用到了时间类动态取值逻辑,提前将时间值存入静态变量再传入计算逻辑,不要直接在DataFrame转换中使用实时取时间的函数;确认任务执行期间上游依赖的数据源不会被其他任务修改。
内容的提问来源于stack exchange,提问作者Yehor Reznichenko
相关产品推荐
相关产品推荐

