Best practices Airflow to create BigQuery tables from table

Viewed 715

I am new to BigQuery and come from an AWS background. I have a bucket with no structure, just files of names YYYY-MM-DD-<SOME_ID>.csv.gzip.

The goal is to import this into BigQuery, then create another dataset with a subset table of the imported data. It should be last month's data, exclude some rows with a WHERE statement and exclude some columns.

There seem to be many alternatives using different operators. What would be the best practice to do it?


BigQueryCreateEmptyDatasetOperator(...)
BigQueryCreateEmptyTableOperator(...)
BigQueryExecuteQueryOperator(...) / BigQueryInsertJobOperator / BigQueryUpsertTableOperator

I also found

from airflow.providers.google.cloud.transfers.gcs_to_bigquery import (
    GCSToBigQueryOperator,
)
GCSToBigQueryOperator(...)

When is this preferred?

This is my current code:

create_new_dataset_A = BigQueryCreateEmptyDatasetOperator(
    dataset_id=DATASET_NAME_A,
    project_id=PROJECT_ID,
    gcp_conn_id='_my_gcp_conn_',
    task_id='create_new_dataset_A')

load_csv = GCSToBigQueryOperator(
    bucket='cloud-samples-data',
    compression="GZIP",
    create_disposition="CREATE_IF_NEEDED",
    destination_project_dataset_table=f"{PROJECT_ID}.{DATASET_NAME_A}.{TABLE_NAME}",
    source_format="CSV",
    source_objects=['202*'],
    task_id='load_csv',
    write_disposition='WRITE_APPEND',
    schema_fields=[
        {'name': 'name', 'type': 'STRING', 'mode': 'NULLABLE'},
        {'name': 'post_abbr', 'type': 'STRING', 'mode': 'NULLABLE'},
    ],
)

create_new_dataset_B = BigQueryCreateEmptyDatasetOperator(
    dataset_id=DATASET_NAME_B,
    project_id=PROJECT_ID,
    gcp_conn_id='_my_gcp_conn_',
    task_id='create_new_dataset_B')

populate_new_dataset_B = BigQueryExecuteQueryOperator(...) / BigQueryInsertJobOperator / BigQueryUpsertTableOperator

Alternatives below:

    populate_new_dataset_B = BigQueryExecuteQueryOperator(
        task_id='load_from_table_a_to_table_b',
        use_legacy_sql=False,
        write_disposition='WRITE_APPEND',
        sql=f'''
        INSERT `{PROJECT_ID}.{DATASET_NAME_A}.D_EXCHANGE_RATE` 
        SELECT col_x, col_y #skip som col from table_a
        FROM
        `{PROJECT_ID}.{DATASET_NAME_A}.S_EXCHANGE_RATE`
        WHERE col_x is not null
        '''

Does it keep track of rows it loaded due to write_disposition='WRITE_APPEND'? Does GCSToBigQueryOperator keep track of metadata or load duplicates?

 populate_new_dataset_B = BigQueryInsertJobOperator(
    task_id="load_from_table_a_to_table_b",
    configuration={
        "query": {
            "query": "{% include 'sql-file.sql' %}",
            "use_legacy_sql": False,
        }
    },
    dag=dag,
)

Is this more for scheduled ETL jobs? Example: https://github.com/simonbreton/Capstone-project/blob/a6563576fa63b248a24d4a1bba70af10f527f6b4/airflow/dags/sql/fact_query.sql. Here they do not use write_disposition='WRITE_APPEND' they use a where statement instead. Why? When to prefer?

Last operator I dont get, when to use it? https://airflow.apache.org/docs/apache-airflow-providers-google/stable/operators/cloud/bigquery.html#howto-operator-bigqueryupserttableoperator

Which operator to use for populate_new_dataset_B?

Appreciate all help.

0 Answers
Related