Great Expectations验证器用SparkDFDataset报错,求解决方案
解决Great Expectations中SparkDFDataset无persist属性的错误及Hive表校验HTML报告生成方案
错误原因
你遇到的AttributeError: 'SparkDFDataset' object has no attribute 'persist'是因为你将Great Expectations封装后的SparkDFDataset对象传入了Spark数据源的batch请求,而GE的Spark执行引擎期望接收的是原生Spark DataFrame——persist()是原生Spark DataFrame的方法,SparkDFDataset并没有这个方法。
直接修复方案(不初始化GE项目)
去掉对原生Spark DataFrame的封装,直接使用从Hive读取的原生DF构建batch请求,修正后的代码如下:
import great_expectations as ge from pyspark.sql import SparkSession # 初始化SparkSession(确保已配置Hive支持) sk = SparkSession.builder.appName("GE_TEST").enableHiveSupport().getOrCreate() sk.sql("use DB1") # 读取Hive表,得到原生Spark DataFrame hive_table = sk.sql("SELECT * FROM TABLEX") # 创建GE上下文与Spark数据源 context = ge.get_context() datasource = context.sources.add_spark("my_spark_datasource") data_asset = datasource.add_dataframe_asset(name="my_df_asset") # 直接传入原生Spark DataFrame构建batch请求 my_batch_request = data_asset.build_batch_request(dataframe=hive_table) # 创建期望套件与Validator expectation_suite_name = "test_hive_table_suite" context.add_or_update_expectation_suite(expectation_suite_name=expectation_suite_name) validator = context.get_validator( batch_request=my_batch_request, expectation_suite_name=expectation_suite_name ) # 添加校验规则示例 validator.expect_column_values_to_not_be_null(column="id") validator.expect_column_values_to_be_between(column="age", min_value=0, max_value=120) # 保存期望套件 validator.save_expectation_suite(discard_failed_expectations=False) # 运行校验并生成HTML报告 checkpoint = context.add_or_update_checkpoint( name="test_hive_checkpoint", validator=validator, ) checkpoint_result = checkpoint.run() context.build_data_docs()
执行完context.build_data_docs()后,GE会在默认路径(未初始化项目时为./great_expectations/uncommitted/data_docs/local_site/)生成HTML报告,打开其中的index.html即可查看校验结果。
更规范的生产级方案(初始化GE项目)
如果需要长期维护校验规则,建议初始化GE项目,直接配置Hive数据源,无需手动读取DataFrame:
- 执行
great_expectations init初始化项目 - 在
great_expectations.yml中配置Spark/Hive数据源(确保SparkSession已配置Hive支持) - 直接从Hive表创建数据资产,构建batch请求
- 后续的校验规则维护、报告生成流程更规范,适合团队协作
避免内存问题的注意事项
- 绝对不要使用
toPandas()将大数据量Spark DataFrame转成Pandas DF,这会把分布式数据拉到本地节点,必然引发内存溢出 - 所有校验逻辑都基于Spark执行引擎完成,GE会自动利用Spark的分布式计算能力处理400万行的数据集
- 若Spark执行时仍有内存问题,可调整Spark的executor内存参数(如
--executor-memory 8g、--driver-memory 4g)
内容的提问来源于stack exchange,提问作者user7579065
相关产品推荐
相关产品推荐

