Airflow 2.1.4中MySQL表归档至BigQuery的DAG问题求助
问题与解决方案
问题背景
当前使用Airflow 2.1.4,需将MySQL表testing_monitor_archive中超过2周的数据归档至BigQuery表monitoring_table。原DAG因误用GCSCreateBucketOperator(仅用于创建GCS桶,无数据上传能力)并传入已废弃的GoogleCloudStorageBucketOperator参数(provide_context、python_callable等)导致报错,需修正流程实现:MySQL取数→转CSV上传GCS→导入BigQuery。
修正方案
- 拆分桶创建与数据上传任务:用
GCSCreateBucketOperator仅负责确保GCS桶存在(添加exists_ok=True避免重复创建报错),数据上传改用PythonOperator结合GCS客户端库实现。 - 处理查询结果转CSV:通过XCom获取
MySqlOperator的查询结果,用Pandas将数据转为CSV格式后上传至GCS。 - 完善BigQuery导入参数:在
GCSToBigQueryOperator中指定CSV格式、表头处理、自动检测schema等必要参数,确保数据正确导入。
修改后的完整代码
from datetime import timedelta, datetime from airflow import DAG from airflow.operators.mysql_operator import MySqlOperator from airflow.operators.python import PythonOperator from airflow.providers.google.cloud.operators.gcs import GCSCreateBucketOperator from airflow.providers.google.cloud.transfers.gcs_to_bigquery import GCSToBigQueryOperator from airflow.utils.dates import days_ago import pandas as pd from google.cloud import storage from io import StringIO # Default arguments for the DAG default_args = { 'owner': 'airflow', 'start_date': days_ago(2), 'depends_on_past': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } # Create a new DAG dag = DAG( 'my_data_pipeline', default_args=default_args, schedule_interval='0 0 * * *', # 每日午夜执行 ) # SQL query to fetch data older than 2 weeks query = """ SELECT * FROM testing_monitor_archive WHERE CREATE_TS < DATE_SUB(NOW(), INTERVAL 2 WEEK) """ # 1. 从MySQL查询数据 fetch_data = MySqlOperator( task_id='fetch_data', mysql_conn_id='my_mysql_connection', sql=query, dag=dag, ) # GCS配置 bucket = 'mysql-archive-gcs-bucket' object_name = 'data/{{ execution_date }}.csv' # 2. 创建GCS桶(若不存在) create_gcs_bucket = GCSCreateBucketOperator( task_id='create_gcs_bucket', bucket_name=bucket, gcp_conn_id='gcp_conn_id', exists_ok=True, # 桶已存在时不报错 dag=dag, ) # 3. 将查询结果转CSV并上传至GCS def upload_data_to_gcs(**kwargs): ti = kwargs['ti'] # 从XCom获取MySQL查询结果 query_results = ti.xcom_pull(task_ids="fetch_data") if not query_results: print("无符合条件的归档数据,跳过上传") return # 将结果转为DataFrame,生成CSV df = pd.DataFrame(query_results) csv_buffer = StringIO() df.to_csv(csv_buffer, index=False, header=True) # 上传至GCS client = storage.Client() bucket_obj = client.get_bucket(bucket) blob = bucket_obj.blob(object_name.format(execution_date=kwargs['execution_date'])) blob.upload_from_string(csv_buffer.getvalue(), content_type='text/csv') upload_to_gcs = PythonOperator( task_id='upload_to_gcs', python_callable=upload_data_to_gcs, provide_context=True, dag=dag, ) # BigQuery配置 dataset_id = 'archive_dataset' table_id = 'monitoring_table' # 4. 从GCS导入数据至BigQuery load_to_bigquery = GCSToBigQueryOperator( task_id='load_to_bigquery', bucket=bucket, source_objects=[object_name], destination_project_dataset_table=f"{dataset_id}.{table_id}", gcp_conn_id='gcp_conn_id', source_format='CSV', skip_leading_rows=1, # 跳过CSV表头 autodetect=True, # 自动检测表结构 write_disposition='WRITE_APPEND', # 追加数据到现有表 dag=dag, ) # 设置任务依赖 create_gcs_bucket >> fetch_data >> upload_to_gcs >> load_to_bigquery
关键说明
GCSCreateBucketOperator仅负责初始化桶,添加exists_ok=True避免重复创建报错。PythonOperator中通过Pandas处理数据转CSV,使用GCS客户端库直接上传,无需本地文件存储。GCSToBigQueryOperator补充了source_format、skip_leading_rows、autodetect等参数,确保CSV数据正确匹配BigQuery表结构,write_disposition设置为追加模式避免覆盖现有数据。
内容的提问来源于stack exchange,提问作者Tushaar
相关产品推荐
相关产品推荐

