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
相关产品推荐
相关产品推荐

