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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 19:20:08