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
相关产品推荐
相关产品推荐

