You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确克隆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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 23:31:00