【问题标题】:Airflow BigQueryOperator: how to save query result in a partitioned Table?Airflow BigQueryOperator:如何将查询结果保存在分区表中?
【发布时间】:2018-11-03 10:17:30
【问题描述】:

我有一个简单的 DAG

from airflow import DAG
from airflow.contrib.operators.bigquery_operator import BigQueryOperator

with DAG(dag_id='my_dags.my_dag') as dag:

    start = DummyOperator(task_id='start')

    end = DummyOperator(task_id='end')
    sql = """
             SELECT *
             FROM 'another_dataset.another_table'
          """
    bq_query = BigQueryOperator(bql=sql,
                            destination_dataset_table='my_dataset.my_table20180524'),
                            task_id='bq_query',
                            bigquery_conn_id='my_bq_connection',
                            use_legacy_sql=False,
                            write_disposition='WRITE_TRUNCATE',
                            create_disposition='CREATE_IF_NEEDED',
                            query_params={})
    start >> bq_query >> end

执行bq_query 任务时,SQL 查询将保存在分片表中。我希望它保存在每日分区表中。为此,我只将destination_dataset_table 更改为my_dataset.my_table$20180524。执行bq_task时出现以下错误:

Partitioning specification must be provided in order to create partitioned table

如何指定 BigQuery 将查询结果保存到每日分区表?我的第一个猜测是在BigQueryOperator 中使用query_params 但我没有找到任何关于如何使用该参数的示例。

编辑:

我正在使用google-cloud==0.27.0 python 客户端...它是 Prod 中使用的客户端 :(

【问题讨论】:

  • CREATE TABLE... PARTITION BY ... AS SELECT... 不工作吗?
  • @ElliottBrossard 我认为它不会起作用,因为 DAG 每天都会执行。使用CREATE ... 将在每次执行后创建表。我只想创建一个新分区而不是整个表。

标签: google-bigquery airflow


【解决方案1】:

您首先需要创建一个 Empty 分区目标表。按照此处的说明:link 创建一个空的分区表

然后再次在气流管道下方运行。 你可以试试代码:

import datetime
from airflow import DAG
from airflow.contrib.operators.bigquery_operator import BigQueryOperator
today_date = datetime.datetime.now().strftime("%Y%m%d")
table_name = 'my_dataset.my_table' + '$' + today_date
with DAG(dag_id='my_dags.my_dag') as dag:
    start = DummyOperator(task_id='start')
    end = DummyOperator(task_id='end')
    sql = """
         SELECT *
         FROM 'another_dataset.another_table'
          """
    bq_query = BigQueryOperator(bql=sql,
                        destination_dataset_table={{ params.t_name }}),
                        task_id='bq_query',
                        bigquery_conn_id='my_bq_connection',
                        use_legacy_sql=False,
                        write_disposition='WRITE_TRUNCATE',
                        create_disposition='CREATE_IF_NEEDED',
                        query_params={'t_name': table_name},
                        dag=dag
                        )
start >> bq_query >> end

所以我所做的是创建了一个动态表名变量并传递给 BQ 运算符。

【讨论】:

  • 我不再收到错误消息:the Task exited with return code 0。但是,没有创建结果表,我在 BigQuery 中看不到它。奇怪
  • 您首先需要创建一个 Empty 分区目标表。按照此处的说明:https://cloud.google.com/bigquery/docs/creating-column-partitions#creating_an_empty_partitioned_table_with_a_schema_definition 创建一个空的分区表
  • 或者你可以使用BigQueryCreateEmptyTableOperator()通过设置time_partitioning参数来创建一个带分区的空表
  • 问题是我不想指定架构。我不认为在没有任何架构的情况下使用 BigQueryCreateEmptyTableOperator() 是允许的。与在 BigQuery 用户界面中一样,创建一个空的每日分区表并且没有任何字段会导致错误
  • schema in BigQueryCreateEmptyTableOperator() 是可选参数。检查this,我们构建了这个运算符,记住我们不想指定任何模式。如果您检查 BigQueryCreateEmptyTableOperator 类的文档字符串 --> Creates a new, empty table in the specified BigQuery dataset, optionally with schema.
【解决方案2】:

这里的主要问题是我无法访问新版本的谷歌云 python API,产品使用版本0.27.0。 所以,为了完成工作,我做了一些坏事和肮脏的事情:

  • 将查询结果保存在分表中,设为table_sharded
  • 得到table_sharded的架构,让它成为table_schema
  • " SELECT * FROM dataset.table_sharded" 查询保存到提供table_schema 的分区表中

所有这些都被抽象为一个使用钩子的运算符。该钩子负责创建/删除表/分区、获取表架构以及在 BigQuery 上运行查询。

看看code。如果有其他解决方案,请告诉我。

【讨论】:

  • 您可以使用(免费)bq 复制表命令代替查询吗?
【解决方案3】:

使用 BigQueryOperator,您可以传递 time_partitioning 参数,该参数将创建提取时间分区表

bq_cmd = BigQueryOperator (
            task_id=                    "task_id",
            sql=                        [query],
            destination_dataset_table=  destination_tbl,
            use_legacy_sql=             False,
            write_disposition=          'WRITE_TRUNCATE',
            time_partitioning=          {'time_partitioning_type':'DAY'},
            allow_large_results=        True,
            trigger_rule=               'all_success',
            query_params=               query_params,
            dag=                        dag
        )

【讨论】:

    【解决方案4】:
    from datetime import datetime,timedelta
    from airflow import DAG
    from airflow.models import Variable
    from airflow.contrib.operators.bigquery_operator import BigQueryOperator
    from airflow.operators.dummy_operator import DummyOperator
    
    DEFAULT_DAG_ARGS = {
        'owner': 'airflow',
        'depends_on_past': False,
        'retries': 2,
        'retry_delay': timedelta(minutes=10),
        'project_id': Variable.get('gcp_project'),
        'zone': Variable.get('gce_zone'),
        'region': Variable.get('gce_region'),
        'location': Variable.get('gce_zone'),
    }
    
    with DAG(
        'test',
        start_date=datetime(2019, 1, 1),
        schedule_interval=None,
        catchup=False,
        default_args=DEFAULT_DAG_ARGS) as dag:
    
        bq_query = BigQueryOperator(
            task_id='create-partition',
            bql="""SELECT
                    * 
                    FROM
                    `dataset.table_name`""",   -- table from which you want to pull data
            destination_dataset_table='project.dataset.table_name' + '$' + datetime.now().strftime('%Y%m%d'),             -- Auto partitioned table in Bq 
            write_disposition='WRITE_TRUNCATE',
            create_disposition='CREATE_IF_NEEDED',
            use_legacy_sql=False,
        )
    

    我建议在 Airflow 中使用变量并创建所有字段并在 DAG 中使用。 通过上面的代码,将在 Bigquery 表中添加今天日期的分区。

    【讨论】:

      猜你喜欢
      • 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
      相关资源
      最近更新 更多