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

PySpark项目中无配置文件编程式配置Great Expectations问询

问题

我正在把Great Expectations验证框架集成到现有PySpark项目里,官方文档大多用JSON/YAML配置,但我的表架构是用Python类定义的,想把验证规则直接放在这些类里。目前我只会用SparkDFDataset做单个期望验证,但不知道怎么编程式构建期望套件、生成验证报告,试了一段代码还触发了运行时异常,求可行的实现方法。

单个期望验证示例代码

spark = SparkSession.builder.master("local[*]").getOrCreate()

df = spark.createDataFrame([
    Row(x=1, y="foo"),
    Row(x=2, y=None),
])
ds = SparkDFDataset(df)

expectation: ExpectationValidationResult = ds.expect_column_values_to_not_be_null("y")
print(expectation.success)

尝试构建套件的报错代码

ds.append_expectation(ExpectationConfiguration(
    expectation_type="expect_column_values_to_not_be_null",
    kwargs={'column': 'y', 'result_format': 'BASIC'},
))
engine = SparkDFExecutionEngine(
    force_reuse_spark_context=True,
)
validator = Validator(
    execution_engine=engine,
    expectation_suite=ds.get_expectation_suite(),
)
res = validator.validate()

解决方案

1. 正确编程式构建期望套件的步骤

无需混用SparkDFDataset和Validator,直接用Validator绑定PySpark DataFrame的方式更稳妥:

from great_expectations.core import ExpectationSuite
from great_expectations.execution_engine import SparkDFExecutionEngine
from great_expectations.validator import Validator
from pyspark.sql import SparkSession, Row

# 初始化Spark会话
spark = SparkSession.builder.master("local[*]").getOrCreate()

# 创建测试DataFrame
df = spark.createDataFrame([
    Row(x=1, y="foo"),
    Row(x=2, y=None),
])

# 创建空的期望套件
suite = ExpectationSuite(expectation_suite_name="my_pyspark_suite")

# 初始化Validator,关联DataFrame和套件
validator = Validator(
    execution_engine=SparkDFExecutionEngine(spark_df=df),
    expectation_suite=suite
)

# 编程式添加期望规则(API和SparkDFDataset完全一致)
validator.expect_column_values_to_not_be_null("y")
validator.expect_column_values_to_be_between("x", min_value=1, max_value=5)

# 执行验证
validation_result = validator.validate()

# 查看验证结果
print(f"整体验证结果: {validation_result.success}")
for result in validation_result.results:
    print(f"期望规则[{result.expectation_config.expectation_type}]结果: {result.success}")

2. 生成可视化验证报告

直接调用validator的build_data_docs方法即可生成HTML报告:

# 生成报告并自动在浏览器打开
validator.build_data_docs()

报告文件也会保存在项目根目录的uncommitted/data_docs/local_site文件夹中。

3. 将验证规则绑定到Python表架构类

可以把验证逻辑封装到你的表架构类中,实现验证与表结构的绑定:

class MyTableSchema:
    def __init__(self, spark_df):
        self.df = spark_df
        self.suite = ExpectationSuite(expectation_suite_name="my_table_suite")
        self.validator = Validator(
            execution_engine=SparkDFExecutionEngine(spark_df=self.df),
            expectation_suite=self.suite
        )
        self._define_expectations()

    def _define_expectations(self):
        # 在这里定义当前表的所有验证规则
        self.validator.expect_column_values_to_not_be_null("y")
        self.validator.expect_column_values_to_be_in_set("y", ["foo", "bar"])
        self.validator.expect_column_values_to_be_between("x", 1, 10)

    def validate(self):
        return self.validator.validate()

# 使用示例
df = spark.createDataFrame([Row(x=1, y="foo"), Row(x=2, y=None)])
table_validator = MyTableSchema(df)
result = table_validator.validate()
print(result.success)

4. 原代码报错原因

你之前的代码错误在于:SparkDFDataset的append_expectation需要配合自身的验证流程,直接将它的套件传给独立Validator会导致执行引擎上下文不匹配,改用Validator直接绑定DataFrame的方式即可避免该问题。

内容的提问来源于stack exchange,提问作者ollik1

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 23:39:16