【发布时间】: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 ...将在每次执行后创建表。我只想创建一个新分区而不是整个表。