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

Python Great Expectations结合Spark如何获取逐行数据校验结果

答案

Great Expectations完全可以实现Spark DataFrame逐行打is_valid标识的需求,默认返回全量列维度统计结果只是默认输出配置的效果,并非工具本身的能力边界。

实现方式

针对Spark引擎的场景,有两种可直接落地的方案:

  • 方案1:自定义行级Expectation直接输出带打标字段的DataFrame
    自定义Expectation时直接把多规则组合的逐行判断逻辑写在_validate方法里,GE会将逻辑下推到Spark原生执行,不会额外带来性能损耗,核心实现参考:
    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供下游使用
                }
            }
    
    跑校验时把Checkpoint的result_format参数设为COMPLETE,就能直接拿到带is_valid字段的结果DataFrame,同时也能拿到全量维度的校验统计指标。
  • 方案2:复用已配置的列级/跨列Expectation规则,自动生成打标逻辑
    如果你已经配置好了所有单字段、字段对的校验规则,不需要重复写判断逻辑:校验运行后可以直接从GE的校验结果中提取每个规则对应的Spark过滤表达式,把所有表达式用AND拼接后,直接在原DataFrame上生成is_valid字段,避免规则双写导致的不一致问题。

注意:执行前必须把GE的执行引擎配置为SparkDFExecutionEngine,不要走Pandas执行路径,否则会把全量数据拉到Driver端计算,引发内存溢出问题;配置正确的前提下,逐行打标的性能和直接手写Spark SQL计算基本无差异。

内容的提问来源于stack exchange,提问作者sunitha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 04:30:50