Delta Live Tables执行Spark SQL报错:无法读取temp3数据集
解决DLT中"Dataset is not defined in the pipeline"错误
你的代码存在几处语法和逻辑问题,导致temp3无法被正确注册,进而触发找不到数据集的错误,以下是具体修正方案:
修正temp2的语法与逻辑错误
原代码错误地将filter操作嵌入表名字符串中,同时存在双引号嵌套的语法问题。正确做法是先读取流表,再调用filter方法过滤数据:@dlt.view def temp2(): return spark.readStream.table("catalog1.schema1.table2").filter(col("field1") == "test")注意:需提前导入
from pyspark.sql.functions import col,否则col函数会报错。修正temp3的SQL拼写错误
SQL语句中把temp2写成了temp 2(多了空格),导致无法识别依赖视图。修正后的代码:@dlt.view def temp3(): return spark.sql("select * from temp1 left join temp2 on temp1.id = temp2.id")验证依赖关系
DLT会自动解析视图间的依赖,但如果上游视图(temp1、temp2)创建失败,temp3也无法被正确注册。修正上述问题后,temp3就能正常生成,后续dlt.read("temp3")即可找到对应数据集。
修正后的完整代码
from pyspark.sql.functions import col @dlt.view def temp1(): return spark.readStream.table("catalog1.schema1.table1") @dlt.view def temp2(): return spark.readStream.table("catalog1.schema1.table2").filter(col("field1") == "test") @dlt.view def temp3(): return spark.sql("select * from temp1 left join temp2 on temp1.id = temp2.id") @dlt.create_table( name="final_table", table_properties={ "quality": "silver", "delta.autoOptimize.optimizeWrite": "true", "delta.autoOptimize.autoCompact": "true", }, ) def populate_final_table(): final_df = dlt.read("temp3") return final_df
内容的提问来源于stack exchange,提问作者SreeVik
相关产品推荐
相关产品推荐

