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

如何让Delta Live Streaming Table仅保留NewData字段?

解决Delta Live Table保留冗余列的问题

问题分析

你遇到的情况是使用select("NewData.*")后,表中仍保留了NewData、OldData原列,同时新增了展开后的字段。这通常和schema定义不匹配或Delta Live Table的schema演化策略有关。

针对性解决方案

情况1:NewData是数组类型(你提到JSON中是数组)

当前schema把NewData定义为StructType(结构体),和实际JSON的数组类型不匹配,导致解析异常,进而出现冗余列。需要修正schema并展开数组:

  1. 修正schema定义,用ArrayType包裹结构体:
from pyspark.sql.types import StructType, StructField, ArrayType, IntegerType, StringType

schema = StructType([
    StructField("NewData", ArrayType(StructType([
        StructField("Field1", IntegerType(), True),
        StructField("Field2", StringType(), True)
        # 补充NewData内的其他字段
    ])), True),
    StructField("OldData", ArrayType(StructType([
        # 补充OldData内的字段
    ])), True)
])
  1. 修改读取逻辑,先展开数组再提取字段:
@dlt.table(
    name="newdata_raw",
    table_properties={"quality": "bronze"},
    temporary=False,
)
def create_table():
    query = (
        spark.readStream.format("cloudFiles")
        .schema(schema)
        .option("cloudFiles.format", "json")
        .option("checkpointLocation", sink_dir +"checkpoint/")
        .load(sink_dir)
        .select(explode("NewData").alias("new_data"))  # 展开NewData数组
        .select("new_data.*")  # 提取数组内对象的所有字段
        .withColumn("load_date", to_timestamp(current_timestamp()))
    )
    return query

情况2:NewData是单个结构体对象

如果JSON中NewData是单个对象而非数组,select("NewData.*")理论上只会保留展开后的字段。出现冗余列是因为旧表已存在,DLT默认的schema演化策略保留了原有列。可以通过以下方式解决:

方案A:重新创建表

直接删除已存在的newdata_raw表,重新运行DLT pipeline,新表将只包含select指定的字段和新增的load_date。

方案B:显式指定字段(避免歧义)

放弃select("NewData.*")的简写,显式列出要保留的字段,确保只提取需要的内容:

@dlt.table(
    name="newdata_raw",
    table_properties={"quality": "bronze"},
    temporary=False,
)
def create_table():
    query = (
        spark.readStream.format("cloudFiles")
        .schema(schema)
        .option("cloudFiles.format", "json")
        .option("checkpointLocation", sink_dir +"checkpoint/")
        .load(sink_dir)
        .select(
            # 逐一列出NewData内的字段
            col("NewData.Field1").alias("Field1"),
            col("NewData.Field2").alias("Field2"),
            # 其他需要的字段
            to_timestamp(current_timestamp()).alias("load_date")
        )
    )
    return query

方案C:关闭schema自动合并

在表属性中添加delta.schema.autoMerge.enabled = false,强制DLT严格遵循当前定义的schema,不再保留旧列:

@dlt.table(
    name="newdata_raw",
    table_properties={
        "quality": "bronze",
        "delta.schema.autoMerge.enabled": "false"
    },
    temporary=False,
)
def create_table():
    # 原有读取逻辑

内容的提问来源于stack exchange,提问作者ggsmith

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 22:59:55