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
相关产品推荐
相关产品推荐

