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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 05:50:21