Python Great Expectations结合Spark如何获取逐行数据校验结果
答案
Great Expectations完全可以实现Spark DataFrame逐行打is_valid标识的需求,默认返回全量列维度统计结果只是默认输出配置的效果,并非工具本身的能力边界。
实现方式
针对Spark引擎的场景,有两种可直接落地的方案:
- 方案1:自定义行级Expectation直接输出带打标字段的DataFrame
自定义Expectation时直接把多规则组合的逐行判断逻辑写在_validate方法里,GE会将逻辑下推到Spark原生执行,不会额外带来性能损耗,核心实现参考:
跑校验时把Checkpoint的from great_expectations.expectations import Expectation import pyspark.sql.functions as F class ExpectRowValidFlag(Expectation): # 可根据实际业务定义入参,比如枚举值范围、字段阈值等 def _validate(self, spark_df): # 组合所有行内校验规则,逐行计算合规标识 tagged_df = spark_df.withColumn( "is_valid", F.when( # 替换为实际业务的多列校验规则,示例如下 F.col("user_id").isNotNull() & F.col("age").between(1, 120) & F.col("pay_amount") >= 0 & F.col("order_status").isin(["paid", "cancelled", "refunded"]), True ).otherwise(False) ) invalid_cnt = tagged_df.filter(F.col("is_valid") == False).count() return { "success": invalid_cnt == 0, "result": { "unexpected_count": invalid_cnt, "tagged_dataframe": tagged_df # 直接返回打标完成的df供下游使用 } }result_format参数设为COMPLETE,就能直接拿到带is_valid字段的结果DataFrame,同时也能拿到全量维度的校验统计指标。 - 方案2:复用已配置的列级/跨列Expectation规则,自动生成打标逻辑
如果你已经配置好了所有单字段、字段对的校验规则,不需要重复写判断逻辑:校验运行后可以直接从GE的校验结果中提取每个规则对应的Spark过滤表达式,把所有表达式用AND拼接后,直接在原DataFrame上生成is_valid字段,避免规则双写导致的不一致问题。
注意:执行前必须把GE的执行引擎配置为
SparkDFExecutionEngine,不要走Pandas执行路径,否则会把全量数据拉到Driver端计算,引发内存溢出问题;配置正确的前提下,逐行打标的性能和直接手写Spark SQL计算基本无差异。
内容的提问来源于stack exchange,提问作者sunitha
相关产品推荐
相关产品推荐

