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

Great Expectations:如何让Validator与Checkpoint处理多文件数据资产?

在Great Expectations中批量处理Spark分布式存储的多个Parquet文件

问题背景

在PySpark特征生成流水线中使用Great Expectations(GE)对中间特征集做数据质量校验,特征集存储在GCS上的2011个.snappy.parquet文件中,已配置InferredAssetGCSDataConnector数据源并成功识别所有文件,但创建BatchRequest后,Validator和Checkpoint仅处理单个文件(返回24行,对应单个文件的行数),需要实现批量处理至少100个文件。

解决方案

1. 通过分区筛选加载指定范围的文件(推荐)

利用数据源配置中解析出的year/month/day分区字段,通过partition_filter筛选指定时间范围的分区,GE会自动加载该范围内的所有文件:

batch_request = BatchRequest(
    datasource_name="Spark_source",
    data_connector_name="dt_partitioned_intermediate_featuresets",
    data_asset_name="intermediate_features",
    batch_spec_passthrough={"reader_method": "parquet"},
    # 示例:筛选2023年1月的所有分区,对应多个文件
    partition_filter={
        "year": "2023",
        "month": "01"
    }
)

该方法利用GCS的分区前缀过滤,无需扫描所有文件,性能最优。

2. 手动指定批量Batch Keys加载任意数量文件

先获取数据资产的所有可用Batch Keys,选择需要的数量(如前100个或随机100个),构建包含多个Batch的请求:

import great_expectations as ge
import random

# 获取GE上下文及数据源组件
context = ge.get_context()
datasource = context.get_datasource("Spark_source")
data_connector = datasource.get_data_connector("dt_partitioned_intermediate_featuresets")

# 获取该数据资产的所有Batch Keys
all_batch_keys = data_connector.get_batch_keys("intermediate_features")
# 选择前100个,或用random.sample(all_batch_keys, 100)随机选
selected_batch_keys = all_batch_keys[:100]

# 构建批量BatchRequest
batch_request = BatchRequest(
    datasource_name="Spark_source",
    data_connector_name="dt_partitioned_intermediate_featuresets",
    data_asset_name="intermediate_features",
    batch_spec_passthrough={"reader_method": "parquet"},
    batch_keys=selected_batch_keys
)

GE会将选中的所有文件合并为一个Spark DataFrame进行校验,符合分布式计算需求。

3. 修改Data Connector配置设置批量大小

在数据源配置中添加batch_size参数,指定每个批次自动加载的文件数量:

intermediate_df_datasource_config = {
    "name": "Spark_source",
    "class_name": "Datasource",
    "execution_engine": {"class_name": "SparkDFExecutionEngine"},
    "data_connectors": {
        "dt_partitioned_intermediate_featuresets": {
            "class_name": "InferredAssetGCSDataConnector",
            "bucket_or_name": "[bucket]",
            "prefix": "[prefix]",
            "default_regex": {
                "pattern": "[prefix](.*)_df/dt=(\\d{4})-(\\d{2})-(\\d{2})(.*)\\.snappy\\.parquet",
                "group_names": ["data_asset_name","year", "month", "day", "partition"],
            },
            # 每个批次加载100个文件
            "batch_size": 100
        }, 
    }
}

重新加载数据源后,创建的BatchRequest会自动按batch_size分组加载文件;若需处理所有文件,可循环遍历所有批次。

注意事项

  • 使用SparkDFExecutionEngine时,多个文件会自动合并为分布式DataFrame,不会单节点加载,适配大规模数据场景。
  • 若需随机选择文件,可结合random.sample()方法实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 13:05:30