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

Azure Databricks DLT中apply_changes的keys参数化失败问题

问题解决:DLT元数据驱动中Keys参数化导致的列解析错误

问题描述

搭建元数据驱动的DLT笔记本,从配置表读取参数将ADLS Gen2的数据处理至DLT表,target、source、sequence_by等参数化均正常,但参数化keys选项时出现异常:传入类似'["pkey"]'的字符串参数后,Spark会将字符串的每个字符当作列名,抛出[UNRESOLVED_COLUMN.WITH_SUGGESTION]错误,提示无法解析[、"等“列名”。

原因分析

配置表中存储的Keys是JSON格式的字符串,但代码中直接将其作为字符串传给dlt.apply_changes的keys参数。而该参数要求接收列名组成的Python列表,直接传入字符串会被当作可迭代对象,逐个字符拆分处理,最终导致列解析错误。

解决方案

使用json.loads()将存储的JSON字符串解析为Python列表,再传给keys参数。需确保配置表中存储的字符串是标准JSON格式(如["pkey"],转义后为'[\"pkey\"]')。

修改后的代码

关键修改点

在读取Keys参数后添加JSON解析代码:

keys = row['Keys']
# 新增:将JSON字符串解析为Python列表
keys = json.loads(keys)

完整修改后的PySpark脚本

import dlt 
from pyspark.sql import SparkSession 
from pyspark.sql.functions import col, lit, expr 
from pyspark.sql import Row 
import json 
from pyspark.sql import functions as F 
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, LongType, TimestampType, DateType, DoubleType

df = spark.sql("SELECT * FROM xyzbank.metadata.sourcemetadata WHERE SourceGroupId = 1")
for row in df.collect():
    schema_query = f"SELECT * FROM xyzbank.metadata.sourceschemaconfig WHERE SourceMetaDataId = {row['SourceMetaDataId']}"
    df_schema = spark.sql(schema_query)
    
    data_type_mapping = {
        "StringType": StringType(),
        "IntegerType": IntegerType(),
        "TimeType": TimestampType(),
        "Datetime": DateType(),
        "DoubleType": DoubleType(),
        "DateType": DateType()
    }

    distinct_datatypes = (
        df_schema.select("ColumnDataType", "ColumnName", "ColumnOrder").distinct().collect()
    )

    distinct_datatypes = sorted(distinct_datatypes, key=lambda x: x.ColumnOrder)

    schema_fields = [
        StructField(row.ColumnName, data_type_mapping[row.ColumnDataType], True)
        for row in distinct_datatypes
    ]

    schema = StructType(schema_fields)
    
    table_name = row['SourceTableName']
    checks = row['SourceDataQuality']
    checks = json.loads(checks)
    keys = row['Keys']
    # 解析JSON字符串为Python列表
    keys = json.loads(keys)
    sequence_by = row['SequenceBy']
    file_path = row['SourceFilePath']
    # 建议替换eval为json.loads,更安全
    cloud_file_options = json.loads(row['SourceFileOptions'])
    dq_rules = "({0})".format("AND".join(checks.values()))

    @dlt.table(
        name = "brz_load_"+table_name
    )
    def bronze_load():
        df3 = spark.readStream.format("cloudFiles").options(**cloud_file_options).schema(schema).load(file_path)
        df3 = df3.withColumn("file_process_date", F.current_timestamp())
        return df3
    
    @dlt.table(
        name = "stag_silver_load_"+table_name
    )
    @dlt.expect_all(checks)
    def stag_silver_table():
        df3 = dlt.readStream("brz_load_"+table_name)
        df3 = df3.withColumn("dq_check", F.expr(dq_rules)).filter("dq_check=true")
        return df3
    

    dlt.create_streaming_table(
        name = "silver_load_"+table_name
    )
    dlt.apply_changes(
        target = "silver_load_"+table_name,
        source = "stag_silver_load_"+table_name,
        keys=keys,
        stored_as_scd_type=2,
        sequence_by=sequence_by
    )

    @dlt.table(
        name = "quarantine_silver_load_"+table_name
    )
    @dlt.expect_all(checks)
    def quarantine_silver_table():
        df3 = dlt.readStream("brz_load_"+table_name)
        df3 = df3.withColumn("dq_check", F.expr(dq_rules)).filter("dq_check=false")
        return df3

配置表数据验证

确保配置表中Keys字段存储标准JSON格式字符串,示例插入语句:

INSERT INTO XyzBank.MetaData.SourceMetaData (SourceFilePath,SourceFileFormat,SourceActive,SourceDelimeter,SourceHeader,SourceDataQuality,SourceFileOptions,SourceTableName,ScdType,SequenceBy,Keys,SourceGroupId,CreatedOn,CreatedBy,ModifiedOn,ModifiedBy)
SELECT 'abfss://abc@xyz.dfs.core.windows.net/DLT/', 'csv', 'True', ',', 'True', '{"validate Description":"(Description is NOT NULL)","validate Name":"(Name is NOT NULL)"}', '{"cloudFiles.format":"csv","header":true}', 'Product', 2,'file_process_date','[\"pkey\"]',1, current_timestamp(), 'ABC', current_timestamp(), 'ABC'

补充说明

  • 建议替换代码中eval(row['SourceFileOptions'])为json.loads(row['SourceFileOptions']),避免eval的安全风险,同时符合JSON解析规范。
  • 多键场景(如'[\"pkey\",\"id\"]')同样适用,解析后会自动生成包含多个列名的列表。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 20:05:09