如何使用PySpark修改托管Delta表的列数据类型,支持基于输入参数操作
PySpark修改Delta表列数据类型实现方案
1. 如何使用PySpark修改托管Delta表的列数据类型
修改托管Delta表的列数据类型可按照「读取表数据→指定列做类型转换→写回表并同步更新schema」的流程完成,写入时需要开启覆写schema的配置,避免Delta的schema校验拦截写入请求。
提示:如果是Delta 2.0及以上版本,针对非分区列的简单类型转换也可以直接执行
ALTER TABLESQL语句修改,无需读取全表数据。
2. 如何基于输入参数使用PySpark调整列数据类型
你可以将待转换的列名、目标数据类型、目标表名都设置为可传入的参数,无需在代码中硬编码对应值,即可实现动态调整列类型的逻辑。
参考实现代码
from pyspark.sql.types import IntegerType, BooleanType, DateType from pyspark.sql.functions import col # 以下为可动态传入的参数 Column_Name = "EFFECTIVE_DATE" Target_Type = DateType() Target_Table = "TableA" # 读取目标表数据 df = spark.sql(f"select * from {Target_Table}") # 对指定列做类型转换 df_converted = df.withColumn(Column_Name, col(Column_Name).cast(Target_Type)) # 转换完成后写回Delta表 df_converted.write.format("delta")\ .mode("overwrite")\ .option("overwriteSchema", "true")\ .saveAsTable(Target_Table)
内容的提问来源于stack exchange,提问作者Venkatesh
相关产品推荐
相关产品推荐

