Airflow 1.9中BigQueryValueCheckOperator如何启用标准SQL?
在Airflow 1.9中让BigQueryValueCheckOperator支持标准SQL的解决方案
好问题!在Apache Airflow 1.9版本中,BigQueryValueCheckOperator确实没有直接提供use_legacy_sql参数来切换到标准SQL——这个operator在当时的实现里默认使用Legacy SQL,而且没有开放这个配置项。不过针对你需要用_PARTITIONTIME的场景,有两种靠谱的解决办法:
方法一:自定义支持标准SQL的Operator
你可以继承原BigQueryValueCheckOperator,添加use_legacy_sql参数并传递给底层的BigQuery Hook,这样就能直接用标准SQL写检查语句了。
from airflow.contrib.operators.bigquery_operator import BigQueryValueCheckOperator from airflow.contrib.hooks.bigquery_hook import BigQueryHook class BigQueryStandardSQLValueCheckOperator(BigQueryValueCheckOperator): def __init__(self, use_legacy_sql=False, location=None, *args, **kwargs): super().__init__(*args, **kwargs) self.use_legacy_sql = use_legacy_sql self.location = location # 可选,若需指定BQ数据集位置 def get_db_hook(self): return BigQueryHook( bigquery_conn_id=self.bigquery_conn_id, use_legacy_sql=self.use_legacy_sql, location=self.location )
使用这个自定义Operator的示例:
check_partition_data = BigQueryStandardSQLValueCheckOperator( task_id='check_daily_partition_count', sql='SELECT COUNT(*) FROM `your-project.your-dataset.your-table` WHERE _PARTITIONTIME >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY)', pass_value=1000, # 你期望的阈值 bigquery_conn_id='your_bigquery_connection', use_legacy_sql=False, dag=your_dag )
方法二:拆分任务,用BigQueryOperator+ValueCheckOperator
如果不想自定义代码,可以把检查拆成两步:先用BigQueryOperator执行标准SQL查询并将结果推送到XCom,再用ValueCheckOperator验证XCom中的值。
from airflow.contrib.operators.bigquery_operator import BigQueryOperator from airflow.operators.check_operator import ValueCheckOperator # 第一步:执行标准SQL查询,将结果推送到XCom run_bq_query = BigQueryOperator( task_id='fetch_partition_record_count', sql='SELECT COUNT(*) AS record_count FROM `your-project.your-dataset.your-table` WHERE _PARTITIONTIME BETWEEN TIMESTAMP("2024-01-01") AND TIMESTAMP("2024-01-02")', use_legacy_sql=False, bigquery_conn_id='your_bigquery_connection', do_xcom_push=True, dag=your_dag ) # 第二步:验证XCom中的查询结果 verify_record_count = ValueCheckOperator( task_id='verify_partition_count', sql="SELECT '{{ ti.xcom_pull(task_ids=\"fetch_partition_record_count\")['record_count'] }}'", pass_value=1000, conn_id='sqlite_default', # 用Airflow默认的元数据库连接即可,仅用于执行简单查询验证 dag=your_dag ) # 设置任务依赖 run_bq_query >> verify_record_count
两种方法各有优势:第一种适合需要多次复用标准SQL检查的场景,用法和原Operator一致;第二种更轻量化,不需要额外的自定义代码,适合一次性的检查任务。
内容的提问来源于stack exchange,提问作者Jean-Christophe Rodrigue
相关产品推荐
相关产品推荐

