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

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做校验——核心思路是把流拆分成一个个独立的批处理任务,对每个批单独执行校验逻辑。

修正后的核心代码逻辑

  1. 转换为GE兼容的DataFrame:Spark原生DataFrame没有GE的expect_*方法,需先转换成ge.dataframe.SparkDFDataset类型。
  2. 处理校验结果:根据校验成功/失败的结果,决定是否写入目标表,或记录错误数据。
  3. 修复变量错误:原代码中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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 19:55:00