如何正确克隆Delta表?Databricks SQL编辑器报列不存在异常
一、解决当前代码报错的方案
1. 替换低效的跨连接数据读取方式
你当前通过databricks sql客户端拉取全表到Pandas再转Spark的流程,不仅效率低下,还容易触发数据类型不兼容、内存溢出等问题。直接在Databricks环境中用Spark读取生产表,跳过中间转换环节:
from pyspark.sql import functions as F # 直接读取生产表 prod_df = spark.read.table("schema._production_table") # 按需调整数据类型(避免手动类型转换的兼容性问题) prod_df = prod_df.withColumn("mandt", F.col("mandt").cast("int")) \ .withColumn("umrez", F.col("umrez").cast("int")) \ .withColumn("atinn", F.col("atinn").cast("int")) \ .withColumn("mesrt", F.col("mesrt").cast("int")) \ .withColumn("nest_ftr", F.col("nest_ftr").cast("int")) \ .withColumn("max_stack", F.col("max_stack").cast("int")) \ .withColumn("__insert_gmt_ts", F.to_timestamp(F.col("__insert_gmt_ts"))) \ .withColumn("__update_gmt_ts", F.to_timestamp(F.col("__update_gmt_ts"))) # 写入测试Delta表 prod_df.write.option("overwriteSchema", "true") \ .mode("overwrite") \ .format("delta") \ .saveAsTable("schema.test_table")
2. 修复现有代码的类型转换问题
如果必须保留原流程,修正以下关键问题:
- 把类型字典中的
str改为object:Pandas的str类型会强制非空,原数据的NULL值会被转为字符串'None',导致Spark读取时类型异常; - 整数列改用Pandas的
Int64(支持空值的整数类型),避免空值被自动转为float类型; - 日期列不要在Pandas中强制转
datetime64[ns],交给Spark自动推断或用to_timestamp转换。
修正后的类型字典:
column_names = { 'mandt':'Int64', 'matnr':object, 'meinh':object, 'umrez':'Int64', 'umren':object, 'eannr':object, 'ean11':object, 'numtp':object, 'laeng':object, 'breit':object, 'hoehe':object, 'meabm':object, 'volum':object, 'voleh':object, 'brgew':object, 'gewei':object, 'mesub':object, 'atinn':'Int64', 'mesrt':'Int64', 'xfhdw':object, 'xbeww':object, 'kzwso':object, 'msehi':object, 'bflme_marm':object, 'gtin_variant':object, 'nest_ftr':'Int64', 'max_stack':'Int64', 'capause':object, 'ty2tq':object, '__insert_gmt_ts':object, '__update_gmt_ts':object }
同时,替换fetchall()为分块读取,避免全量数据拉取导致的内存溢出:
chunk_size = 10000 data = [] while True: chunk = cursor.fetchmany(chunk_size) if not chunk: break data.extend(chunk) data = pd.DataFrame(data, columns=column_names)
二、基于SQL查询的Delta表克隆指南
Delta Lake原生的CLONE命令仅支持全表(含版本、历史)克隆,若要基于SQL查询结果克隆(复制子集或特定列),可使用以下两种方式:
1. CTAS(Create Table As Select)直接创建Delta表
通过SQL语句直接将查询结果写入新的Delta表,是最简洁的方式:
CREATE OR REPLACE TABLE schema.test_table USING delta AS SELECT mandt, matnr, meinh, umrez, umren, eannr, __insert_gmt_ts, __update_gmt_ts FROM schema._production_table WHERE mandt = '100' -- 可选:添加过滤条件
如果需要保留原表的Schema和约束,可以先克隆空表再插入数据:
-- 克隆原表Schema但不复制数据 CREATE TABLE schema.test_table CLONE schema._production_table WHERE 1=0; -- 插入查询结果 INSERT OVERWRITE schema.test_table SELECT * FROM schema._production_table WHERE mesrt > 0;
2. Spark SQL动态写入
在Notebook中通过Spark SQL执行插入操作,适合动态调整查询条件的场景:
spark.sql(""" INSERT OVERWRITE schema.test_table SELECT * FROM schema._production_table WHERE mesrt > 0 """)
内容的提问来源于stack exchange,提问作者Xkid
相关产品推荐
相关产品推荐

