如何为Great Expectations数据集强制添加BigQuery分区过滤器?
我在BigQuery有两张按DATE类型列_date分区的事件表:
- tableA无需分区过滤器即可查询
- tableB必须添加
_date过滤器才能查询(无法修改表结构或这个强制要求)
在配置Great Expectations数据源时,已经指定按_date列拆分批次,但tableB仍报错提示未提供过滤器。我想通过ConfiguredAssetSqlDataConnector强制添加该过滤器,但官方文档里没找到方法。
当前数据源配置代码
table_name = "the_dataset.tableA" bq_table_datasource_config = { "name": f"{table_name}_bigquery_datasource", "class_name": "Datasource", "module_name": "great_expectations.datasource", "execution_engine": { "module_name": "great_expectations.execution_engine", "class_name": "SqlAlchemyExecutionEngine", "connection_string": "bigquery://our-bigquery-project", }, "data_connectors": { "daily": { "class_name": "ConfiguredAssetSqlDataConnector", "include_schema_name": True, "assets": { f"our-bigquery-project.{table_name}": { "splitter_method": "split_on_year_and_month_and_day", "splitter_kwargs": { "column_name": "_date" }, } }, }, }, }
测试代码
context.test_yaml_config( yaml.dump(bq_table_datasource_config), return_mode="report_object", shorten_tracebacks=True, )
报错信息
Attempting to instantiate class from config...
Instantiating as a Datasource, since class_name is Datasource
Successfully instantiated DatasourceExecutionEngine class name: SqlAlchemyExecutionEngine
Data Connectors:
Traceback (most recent call last):
File "/usr/local/lib/python3.9/site-packages/google/cloud/bigquery/dbapi/cursor.py", line 203, in _execute
self._query_job.result()
google.api_core.exceptions.BadRequest: 400 Cannot query over table 'our-bigquery-project.the_dataset.tableB' without a filter over column(s) '_date' that can be used for partition eliminationLocation: US
Job ID: 00516791-f350-4f82-a8a1-d2e60bc08ec5During handling of the above exception, another exception occurred:
Traceback (most recent call last):
File "/usr/local/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 1799, in _execute_context
self.dialect.do_execute(
google.cloud.bigquery.dbapi.exceptions.DatabaseError: 400 Cannot query over table 'our-bigquery-project.the_dataset.tableB' without a filter over column(s) '_date' that can be used for partition eliminationLocation: US
Job ID: 00000000-ffff-4444-aaaa-121212121212The above exception was the direct cause of the following exception:
Traceback (most recent call last):
File "/usr/local/lib/python3.9/site-packages/great_expectations/data_context/config_validator/yaml_config_validator.py", line 266, in test_yaml_config
report_object: dict = instantiated_class.self_check(
sqlalchemy.exc.DatabaseError: (google.cloud.bigquery.dbapi.exceptions.DatabaseError) 400 Cannot query over table 'our-bigquery-project.the_dataset.tableB' without a filter over column(s) '_date' that can be used for partition eliminationLocation: US
Job ID: 00000000-ffff-4444-aaaa-121212121212[SQL: SELECT distinct(concat(concat(concat(%(concat_1:STRING)s, CAST(EXTRACT(year FROM
_date) AS STRING)), CAST(EXTRACT(month FROM_date) AS STRING)), CAST(EXTRACT(day FROM_date) AS STRING))) ASconcat_distinct_values, CAST(EXTRACT(year FROM_date) AS INT64) ASyear, CAST(EXTRACT(month FROM_date) AS INT64) ASmonth, CAST(EXTRACT(day FROM_date) AS INT64) ASday
FROM our-bigquery-project.the_dataset.tableB]
[parameters: {'concat_1': ''}]
解决方案
方法1:添加batch_filter_parameters指定默认过滤条件
在tableB的资产配置中添加batch_filter_parameters,强制每次查询都带上_date过滤条件,比如限定最近30天的数据:
table_name = "the_dataset.tableB" bq_table_datasource_config = { "name": f"{table_name}_bigquery_datasource", "class_name": "Datasource", "module_name": "great_expectations.datasource", "execution_engine": { "module_name": "great_expectations.execution_engine", "class_name": "SqlAlchemyExecutionEngine", "connection_string": "bigquery://our-bigquery-project", }, "data_connectors": { "daily": { "class_name": "ConfiguredAssetSqlDataConnector", "include_schema_name": True, "assets": { f"our-bigquery-project.{table_name}": { "splitter_method": "split_on_year_and_month_and_day", "splitter_kwargs": { "column_name": "_date" }, # 强制添加分区过滤条件 "batch_filter_parameters": { "_date": "> DATE_SUB(CURRENT_DATE(), INTERVAL 30 DAY)" } } }, }, }, }
方法2:使用自定义SQL查询作为资产
如果需要更灵活的过滤逻辑,可以将资产定义为带过滤条件的自定义SQL查询,替代直接引用表:
table_name = "the_dataset.tableB" bq_table_datasource_config = { "name": f"{table_name}_bigquery_datasource", "class_name": "Datasource", "module_name": "great_expectations.datasource", "execution_engine": { "module_name": "great_expectations.execution_engine", "class_name": "SqlAlchemyExecutionEngine", "connection_string": "bigquery://our-bigquery-project", }, "data_connectors": { "daily": { "class_name": "ConfiguredAssetSqlDataConnector", "include_schema_name": True, "assets": { f"our-bigquery-project.{table_name}": { "splitter_method": "split_on_year_and_month_and_day", "splitter_kwargs": { "column_name": "_date" }, # 用带过滤的SQL查询替代直接表引用 "query": "SELECT * FROM our-bigquery-project.the_dataset.tableB WHERE _date >= DATE_SUB(CURRENT_DATE(), INTERVAL 30 DAY)" } }, }, }, }
方法3:自定义DataConnector子类(高级场景)
如果以上方法无法满足需求,可以自定义ConfiguredAssetSqlDataConnector的子类,重写生成查询的方法,确保所有自动生成的查询都强制包含_date过滤条件。
内容的提问来源于stack exchange,提问作者dlamblin

