Databricks DLT循环处理表时触发‘Cannot redefine dataset’异常求助
问题
单表处理的代码运行正常,但用for循环处理数据库所有表时,抛出错误:"AnalysisException: Cannot redefine dataset 'source_ds',Map(),Map(),List(),List(),Map())"。需求是将表名传入source_ds,基于键和序列列处理CDC数据。
代码如下:
import dlt from pyspark.sql.functions import * from pyspark.sql.types import * import time raw_db_name = "raw_db" def generate_silver_tables(target_table, source_table, keys_col_list): @dlt.table def source_ds(): return spark.table(f"{raw_db_name}.{source_table}") ### Create the target table definition dlt.create_target_table(name=target_table, comment= f"Clean, merged {target_table}", #partition_cols=["topic"], table_properties={ "quality": "silver", "pipelines.autoOptimize.managed": "true" } ) ## Do the merge dlt.apply_changes( target = target_table, source = "source_ds", keys = keys_col_list, apply_as_deletes = expr("operation = 'DELETE'"), sequence_by = col("ts_ms"), ignore_null_updates = False, except_column_list = ["operation", "timestamp_ms"], stored_as_scd_type = "1" ) return # THIS WORKS FINE #--------------- # raw_dbname = "raw_db" # raw_tbl_name = 'raw_table' # processed_tbl_name = raw_tbl_name.replace("raw", "processed") # generate_silver_tables(processed_tbl_name, raw_tbl_name) table_list = spark.sql(f"show tables in landing_db ").collect() for row in table_list: landing_tbl_name = row.tableName s2 = spark.sql(f"select key from {landing_db_name}.{landing_tbl_name} limit 1") keys_col_list = list(json.loads(s2.collect()[0][0]).keys()) raw_tbl_name = landing_tbl_name.replace("landing", "raw") processed_tbl_name = landing_tbl_name.replace("landing", "processed") generate_silver_tables(processed_tbl_name, raw_tbl_name, keys_col_list) # time.sleep(10)
解决方案
错误根源是每次调用generate_silver_tables时,都定义了同名的source_ds DLT表,DLT不允许重复注册同名数据集。以下两种方案可解决问题:
方案1:动态生成唯一的source数据集名称
为每个表生成专属的source_ds名称,避免重复定义:
import dlt from pyspark.sql.functions import * from pyspark.sql.types import * import json # 原代码遗漏的导入 raw_db_name = "raw_db" landing_db_name = "landing_db" # 原代码未定义的变量 def generate_silver_tables(target_table, source_table, keys_col_list): # 结合表名生成唯一的source数据集名称 source_ds_name = f"source_ds_{source_table}" @dlt.table(name=source_ds_name) def source_ds(): return spark.table(f"{raw_db_name}.{source_table}") # 创建目标表 dlt.create_target_table( name=target_table, comment=f"Clean, merged {target_table}", table_properties={ "quality": "silver", "pipelines.autoOptimize.managed": "true" } ) # 应用CDC变更时使用动态生成的source名称 dlt.apply_changes( target=target_table, source=source_ds_name, keys=keys_col_list, apply_as_deletes=expr("operation = 'DELETE'"), sequence_by=col("ts_ms"), ignore_null_updates=False, except_column_list=["operation", "timestamp_ms"], stored_as_scd_type="1" ) # 遍历处理所有表 table_list = spark.sql(f"show tables in {landing_db_name}").collect() for row in table_list: landing_tbl_name = row.tableName s2 = spark.sql(f"select key from {landing_db_name}.{landing_tbl_name} limit 1") keys_col_list = list(json.loads(s2.collect()[0][0]).keys()) raw_tbl_name = landing_tbl_name.replace("landing", "raw") processed_tbl_name = landing_tbl_name.replace("landing", "processed") generate_silver_tables(processed_tbl_name, raw_tbl_name, keys_col_list)
方案2:用临时视图替代命名DLT表
跳过定义source_ds DLT表,直接使用临时视图传递数据源,更轻量化:
import dlt from pyspark.sql.functions import * from pyspark.sql.types import * import json raw_db_name = "raw_db" landing_db_name = "landing_db" def generate_silver_tables(target_table, source_table, keys_col_list): # 创建临时视图 spark.table(f"{raw_db_name}.{source_table}").createOrReplaceTempView(f"temp_source_{source_table}") # 创建目标表 dlt.create_target_table( name=target_table, comment=f"Clean, merged {target_table}", table_properties={ "quality": "silver", "pipelines.autoOptimize.managed": "true" } ) # 应用CDC变更时引用临时视图 dlt.apply_changes( target=target_table, source=f"temp_source_{source_table}", keys=keys_col_list, apply_as_deletes=expr("operation = 'DELETE'"), sequence_by=col("ts_ms"), ignore_null_updates=False, except_column_list=["operation", "timestamp_ms"], stored_as_scd_type="1" ) # 遍历处理所有表 table_list = spark.sql(f"show tables in {landing_db_name}").collect() for row in table_list: landing_tbl_name = row.tableName s2 = spark.sql(f"select key from {landing_db_name}.{landing_tbl_name} limit 1") keys_col_list = list(json.loads(s2.collect()[0][0]).keys()) raw_tbl_name = landing_tbl_name.replace("landing", "raw") processed_tbl_name = landing_tbl_name.replace("landing", "processed") generate_silver_tables(processed_tbl_name, raw_tbl_name, keys_col_list)
额外说明
- 原代码遗漏了
import json和landing_db_name变量定义,已在修正后的代码中补充 - 方案1会生成多个中间DLT表,若不需要可在Pipeline配置中添加清理规则;方案2的临时视图不会留存,更适合批量处理场景
内容的提问来源于stack exchange,提问作者Yuva
相关产品推荐
相关产品推荐

