Great Expectations结合Spark如何获取校验不通过的完整异常行
问题根因
这是SparkDFExecutionEngine的默认设计行为,和PandasExecutionEngine的逻辑存在差异:Spark引擎出于减少shuffle数据量、控制校验作业执行开销的考量,默认只会返回不符合校验规则的异常字段值本身,即便传入include_unexpected_rows=True,也不会主动拉取异常值对应行的其余字段,不属于参数传错的问题。
可行解决方案
方案1:开启Spark引擎全异常行返回配置(适用于GE 0.15.0及以上版本)
不需要手动重读数据集做过滤,只需要在初始化Validator时显式打开Spark引擎的整行返回开关,搭配原有校验参数即可直接拿到完整异常行,示例代码如下:
from great_expectations.data_context import DataContext context = DataContext() validator = context.get_validator( datasource_name="your_spark_datasource", data_connector_name="your_data_connector", data_asset_name="your_target_dataset", # 核心配置段 runtime_configuration={ "spark_config": { "spark.sql.execution.arrow.pyspark.enabled": "true" }, "return_full_unexpected_rows": True } ) # 原有校验逻辑无需修改 validate_result = validator.expect_column_values_to_not_be_null( column="EmailAddress", result_format="COMPLETE", include_unexpected_rows=True ) # 校验结果的result.unexpected_rows字段中会包含所有异常行的全字段内容 full_unexpected_rows = validate_result["result"]["unexpected_rows"]
配置说明:
- 开启Arrow支持是为了避免Spark行对象序列化到Python进程时出现类型转换错误
- 0.15.0以下的GE版本没有实现Spark引擎的全异常行返回逻辑,该配置不生效,建议升级到0.15.40以上的稳定版本使用
方案2:基于异常索引直接过滤原有DataFrame(适配所有GE版本)
如果暂时无法升级GE版本,可以直接利用校验结果返回的异常行索引,在Validator已经持有的Spark DataFrame上做过滤拿到完整异常行,不需要重新读取全量数据集,性能开销很低,示例代码如下:
from pyspark.sql import functions as F validate_result = validator.expect_column_values_to_not_be_null( column="EmailAddress", result_format="COMPLETE" ) # 拿到异常行的索引列表 unexpected_indexes = validate_result["result"]["unexpected_index_list"] # 直接在已有的DataFrame上过滤,不需要重读数据 full_unexpected_rows = validator.active_dataframe.filter( F.monotonically_increasing_id().isin(unexpected_indexes) ).collect()
使用注意事项
- 当数据集规模较大、异常行量级超过10万条时,不要直接用
result_format="COMPLETE"把全量异常行拉取到驱动节点,容易触发驱动端OOM。这种场景建议搭配GE自带的Validation Action,直接把异常行写入对应的存储表或对象存储路径,不要在校验接口的返回结果中直接携带全量异常数据。
内容的提问来源于stack exchange,提问作者macdrai
相关产品推荐
相关产品推荐

