【问题标题】:Extend BigQueryExecuteQueryOperator with additional labels using jinja2使用 jinja2 使用附加标签扩展 BigQueryExecuteQueryOperator
【发布时间】:2020-02-23 19:59:25
【问题描述】:

为了track GCP costs using labels,想用一些额外的标签扩展BigQueryExecuteQueryOperator,以便每个任务实例在其构造函数中自动获取这些标签。

class ExtendedBigQueryExecuteQueryOperator(BigQueryExecuteQueryOperator):

    @apply_defaults
    def __init__(self,
                 *args,
                 **kwargs) -> None:
        task_labels = {
            'dag_id': '{{ dag.dag_id }}',
            'task_id': kwargs.get('task_id'),
            'ds': '{{ ds }}',
            # ugly, all three params got in diff. ways
        }
        super().__init__(*args, **kwargs)
        if self.labels is None:
            self.labels = task_labels
        else:
            self.labels.update(task_labels)

with DAG(dag_id=...,
         start_date=...,
         schedule_interval=...,
         default_args=...) as dag:

    t1 = ExtendedBigQueryExecuteQueryOperator(
        task_id=f't1',
        sql=f'SELECT 1;',
        labels={'some_additional_label2':'some_additional_label2'}
        # all labels should be: dag_id, task_id, ds, some_additional_label2
    )

    t2 = ExtendedBigQueryExecuteQueryOperator(
        task_id=f't2',
        sql=f'SELECT 2;',
        labels={'some_additional_label3':'some_additional_label3'}
        # all labels should be: dag_id, task_id, ds, some_additional_label3
    )

    t1 >> t2

然后我丢失了任务级别标签some_additional_label2some_additional_label3

【问题讨论】:

  • 我在您的 DAG 中看到一个错字:lables 而不是 t1t2 中的 labels

标签: google-bigquery airflow


【解决方案1】:

您可以在airflow_local_settings.py 中创建以下policy

def policy(task):
    if task.__class__.__name__ == "BigQueryExecuteQueryOperator":
        task.labels.update({'dag_id': task.dag_id, 'task_id': task.task_id})

来自文档:

您的本地 Airflow 设置文件可以定义一个策略函数,该函数能够根据其他任务或 DAG 属性改变任务属性。它接收单个参数作为对任务对象的引用,并有望改变其属性。

更多政策应用详情:https://airflow.readthedocs.io/en/1.10.9/concepts.html#cluster-policy

在这种情况下,您无需扩展 BigQueryExecuteQueryOperator。唯一缺少的部分是您可以在任务本身中设置的 execution_date

例子:

with DAG(dag_id=...,
         start_date=...,
         schedule_interval=...,
         default_args=...) as dag:

    t1 = BigQueryExecuteQueryOperator(
        task_id=f't1',
        sql=f'SELECT 1;',
        lables={'some_additional_label2':'some_additional_label2', 'ds': '{{ ds }}'}
    )

airflow_local_settings 文件需要在您的 PYTHONPATH 上。你可以放在$AIRFLOW_HOME/config 或者你的 dags 目录下。

【讨论】:

  • 我看到从 1.10.4 开始支持此功能。谢谢 试试看!低版本有什么技巧吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多