Databricks Auto Loader流管道集成Great Expectations遇错求解
问题:Great Expectations适配Spark流数据(Auto Loader管道)的方案
我已基于Auto Loader实现了Bronze→Silver→Gold的数据管道,希望通过Great Expectations执行数据质量校验,但运行以下校验代码时遇到错误:
validator.expect_column_values_to_not_be_null(column="col1") validator.expect_column_values_to_be_in_set( column="col2", value_set=[1,6] )
报错信息:
MetricResolutionError: Queries with streaming sources must be executed with writeStream.start();
看起来Great Expectations默认仅支持静态/批处理数据,请问如何让它适配流数据?
我已按照Databricks官方文档在Notebook中完成Great Expectations初始化,管道核心代码如下:
from pyspark.sql.functions import col, to_date, date_format from pyspark.sql.types import StructType, StructField, StringType, IntegerType, FloatType, DateType import time # autoloader table and checkpoint paths basepath = "/mnt/autoloaderdemodl/datagenerator/" bronzeTable = basepath + "bronze/" bronzeCheckpoint = basepath + "checkpoint/bronze/" bronzeSchema = basepath + "schema/bronze/" silverTable = basepath + "silver/" silverCheckpoint = basepath + "checkpoint/silver/" landingZoneLocation = "/mnt/autoloaderdemodl/datageneratorraw/customerdata_csv" # Load data from the CSV file using Auto Loader to bronze table using rescue as schema evolution option raw_df = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("cloudFiles.schemaEvolutionMode", "rescue") \ .option("Header", True) \ .option("cloudFiles.schemaLocation", bronzeSchema) \ .option("cloudFiles.inferSchema", "true") \ .option("cloudFiles.inferColumnTypes", True) \ .load(landingZoneLocation) # Write raw data to the bronze layer bronze_df = raw_df.writeStream.format("delta") \ .trigger(once=True) \ .queryName("bronzeLoader") \ .option("checkpointLocation", bronzeCheckpoint) \ .option("mergeSchema", "true") \ .outputMode("append") \ .start(bronzeTable) # Wait for the bronze stream to finish bronze_df.awaitTermination() bronze = spark.read.format("delta").load(bronzeTable) bronze_count = bronze.count() display(bronze) print("Number of rows in bronze table: {}".format(bronze_count)) bronze_df = spark.readStream.format("delta").load(bronzeTable) # Apply date format transformations to the DataFrame # Transform the date columns silver_df = bronze_df.withColumn("date1", to_date(col("date1"), "yyyyDDD"))\ .withColumn("date2", to_date(col("date2"), "yyyyDDD"))\ .withColumn("date3", to_date(col("date3"), "MMddyy")) # Write the transformed DataFrame to the Silver layer silver_stream = silver_df.writeStream \ .format("delta") \ .outputMode("append") \ .option("mergeSchema", "true") \ .option("checkpointLocation", silverCheckpoint) \ .trigger(once=True) \ .start(silverTable) # Wait for the write stream to complete silver_stream.awaitTermination() # Count the number of rows in the Silver table silver = spark.read.format("delta").load(silverTable) display(silver) silver_count = silver.count() print("Number of rows in silver table: {}".format(silver_count))
注:暂不考虑使用DLT。
我尝试用foreachBatch集成校验逻辑,但仍有问题,代码如下:
import great_expectations as ge from great_expectations.datasource.types import BatchKwargs bronze_df = spark.readStream.format("delta").load(bronzeTable) # Apply date format transformations to the DataFrame # Transform the date columns silver_df = bronze_df.withColumn("date1", to_date(col("date1"), "yyyyDDD"))\ .withColumn("date2", to_date(col("date2"), "yyyyDDD"))\ .withColumn("date3", to_date(col("date3"), "MMddyy")) def validate_micro_batch(batch_df, epoch): print("inside function") # Use Great Expectations to validate the batch DataFrame clean_df = batch_df clean_df.expect_column_values_to_not_be_null(column="col1") clean_df.expect_column_values_to_be_between( column="col2", min_value=0, max_value=1000 ) clean_df.write.format("delta").option("mergeSchema", "true").mode("append").saveAsTable(silverTable) # Print the validation results for the batch validation_results = clean_df.validate() print("Validation results for batch {}:".format(batch_id)) print(validation_results) # Write the transformed DataFrame to the Silver layer if it passes all expectations silver_stream = silver_df.writeStream \ .format("delta") \ .outputMode("append") \ .foreachBatch(validate_micro_batch) \ .option("checkpointLocation", silverCheckpoint) \ .trigger(once=True) \ .start() # Wait for the write stream to complete silver_stream.awaitTermination() # Count the number of rows in the Silver table silver = spark.read.format("delta").load(silverTable) display(silver) silver_count = silver.count() print("Number of rows in silver table: {}".format(silver_count))
解决方案:通过
foreachBatch将Great Expectations校验嵌入流处理 Great Expectations本身不直接支持流数据集,但可以利用Spark Streaming的foreachBatch机制,将每个微批的静态DataFrame传入Great Expectations做校验——核心思路是把流拆分成一个个独立的批处理任务,对每个批单独执行校验逻辑。
修正后的核心代码逻辑
- 转换为GE兼容的DataFrame:Spark原生DataFrame没有GE的
expect_*方法,需先转换成ge.dataframe.SparkDFDataset类型。 - 处理校验结果:根据校验成功/失败的结果,决定是否写入目标表,或记录错误数据。
- 修复变量错误:原代码中
batch_id未定义,需用foreachBatch传入的epoch_id参数替代。
修正后的完整代码:
import great_expectations as ge from great_expectations.data_context import DataContext # 初始化Great Expectations上下文(如果未全局初始化) context = DataContext.create(project_root_dir='/dbfs/great_expectations') bronze_df = spark.readStream.format("delta").load(bronzeTable) # 数据转换逻辑保持不变 silver_df = bronze_df.withColumn("date1", to_date(col("date1"), "yyyyDDD"))\ .withColumn("date2", to_date(col("date2"), "yyyyDDD"))\ .withColumn("date3", to_date(col("date3"), "MMddyy")) def validate_micro_batch(batch_df, epoch_id): # 将Spark DataFrame转换为Great Expectations的SparkDFDataset ge_df = ge.dataset.SparkDFDataset(batch_df) # 定义数据质量校验规则 ge_df.expect_column_values_to_not_be_null(column="col1", result_format="COMPLETE") ge_df.expect_column_values_to_be_between(column="col2", min_value=0, max_value=1000, result_format="COMPLETE") # 执行校验并获取结果 validation_result = ge_df.validate() # 处理校验结果:成功则写入Silver层,失败则记录异常 if validation_result["success"]: batch_df.write.format("delta")\ .option("mergeSchema", "true")\ .mode("append")\ .save(silverTable) print(f"Batch {epoch_id} 校验通过,已写入Silver层") else: # 可选:将错误数据写入隔离表或记录日志 print(f"Batch {epoch_id} 校验失败,错误详情: {validation_result['results']}") # batch_df.write.format("delta").mode("append").save("/path/to/failed-batches") # 启动流处理,绑定校验逻辑 silver_stream = silver_df.writeStream \ .format("delta") \ .outputMode("append") \ .foreachBatch(validate_micro_batch) \ .option("checkpointLocation", silverCheckpoint) \ .trigger(once=True) \ .start() silver_stream.awaitTermination() # 验证结果 silver = spark.read.format("delta").load(silverTable) display(silver) print(f"Number of rows in silver table: {silver.count()}")
关键注意事项
- GE上下文初始化:确保在Notebook中已正确初始化Great Expectations上下文,或在回调函数内显式创建。
- 微批隔离性:每个
foreachBatch传入的batch_df是静态DataFrame,完全兼容GE的所有校验规则。 - 错误数据处理:建议将校验失败的批次写入单独的错误表,便于后续排查,避免直接丢弃数据。
- 性能优化:如果微批数据量极大,可通过调整
maxFilesPerTrigger等参数控制微批大小,平衡校验性能与流处理延迟。
内容的提问来源于stack exchange,提问作者Chhaya Vishwakarma
相关产品推荐
相关产品推荐

