You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow中MySQL转BigQuery的DAG出现无效参数警告问题

Fixing Airflow Deprecation Warnings for MySQL → GCS → BigQuery DAG

Hey there, let's tackle those deprecation warnings popping up in your ETL DAG. They're caused by invalid arguments passed to your operators, and fixing them is straightforward:

1. Fix the typo in MySqlToGoogleCloudStorageOperator

The first warning is just a simple typo—you've got an extra n in the parameter name! The correct argument is google_cloud_storage_conn_id, not google_cloud_storage_connn_id. That's why Airflow is flagging it as invalid.

Corrected extract task snippet:

extract = MySqlToGoogleCloudStorageOperator(
    task_id="extract_mysql_%s_%s"%(connection,table),
    mysql_conn_id = connection,
    google_cloud_storage_conn_id = 'podioGCPConnection',  # Fixed the typo here
    sql = "SELECT *, '%s' as source FROM podiodb.%s"%(connection,table),
    bucket='podio-reader-storage',
    filename= '%s/%s/%s{}.json'%(connection,table,table),
    schema_filename='%s/schemas/%s.json'%(connection,table),
    dag=dag
)

2. Remove the redundant project_id from GoogleCloudStorageToBigQueryOperator

The second warning happens because GoogleCloudStorageToBigQueryOperator doesn't accept a project_id parameter in pre-Airflow 2.0 versions. You already specify the project implicitly in your destination_project_dataset_table value—quick note: BigQuery uses dots instead of slashes for the dataset.table format, so I adjusted that in the snippet below to avoid future issues.

Corrected load task snippet:

load = GoogleCloudStorageToBigQueryOperator(
    task_id = "load_bg_%s_%s"%(connection,table),
    bigquery_conn_id = 'podioGCPConnection',
    google_cloud_storage_conn_id = 'podioGCPConnection',
    bucket = 'podio-reader-storage',
    destination_project_dataset_table = "Podio_Data1.%s.%s"%(connection,table),  # Fixed slash to dot here
    source_objects = ["%s/%s/%s*.json"%(connection,table,table)],
    schema_object = "%s/schemas/%s.json"%(connection,table),
    source_format = 'NEWLINE_DELIMITED_JSON',
    create_disposition = 'CREATE_IF_NEEDED',
    write_disposition = 'WRITE_TRUNCATE',
    # Removed redundant project_id parameter here
    dag=dag
)

Quick Bonus Fix

Right now, your slack_notify.set_upstream(load) line only sets the Slack alert to depend on the last load task created in your loops. If you want the notification to trigger after all load tasks finish, collect all load tasks in a list first:

# Initialize a list to track all load tasks
all_load_tasks = []

for connection in my_connections:
    for table in my_tables:
        # ... extract task code ...
        # ... load task code ...
        load.set_upstream(extract)
        all_load_tasks.append(load)

# Set Slack alert to run after every load task completes
slack_notify.set_upstream(all_load_tasks)

内容的提问来源于stack exchange,提问作者Tia

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 07:04:05