Airflow中MySQL转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

