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

GCP Cloud Composer 2实现MySQL到GCS及BigQuery的Upsert同步问题

MySQL到GCS再到BigQuery的Upsert同步解决方案

1. 替换错误的BigQueryUpsertTableOperator实现Upsert

你误用了BigQueryUpsertTableOperator——这个Operator的作用是更新BigQuery表的元数据(比如修改schema、添加描述),不能用来加载数据或执行Upsert。正确的Upsert流程需要分两步:

  • 将GCS中的CSV数据加载到BigQuery临时表
  • 执行MERGE SQL语句,将临时表的数据合并到目标表(根据主键判断:存在则更新,不存在则插入)

具体实现用GCSToBigQueryOperator加载临时表,再用BigQueryExecuteQueryOperator执行MERGE。

2. 正确实现增量数据拉取

原代码用全局Variable存储执行日期存在并发风险,且用datetime.now()会导致时间漂移,应该基于Airflow的调度时间实现:

  • 利用Airflow模板变量{{ prev_execution_date }}和{{ execution_date }},第一次运行时prev_execution_date为None,拉取全量数据
  • 用参数化SQL避免注入风险
  • 调度时间与数据时间范围严格对齐,保证数据一致性

3. 配置自定义Schema同步数据类型

针对MySQL导出到GCS时的数据类型偏差,可通过以下方式修正:

  • 在MySQL查询时主动转换数据类型(比如将BLOB转为Base64字符串,BigQuery加载后再转成BYTES;日期时间转成ISO标准格式)
  • 在GCSToBigQueryOperator中通过schema_fields指定精确的BigQuery数据类型
  • 提前创建好BigQuery目标表的schema,保证临时表与目标表类型匹配

修正后的完整DAG代码

from datetime import timedelta, datetime
from airflow import DAG
from airflow.providers.google.cloud.transfers.mysql_to_gcs import MySQLToGCSOperator
from airflow.providers.google.cloud.transfers.gcs_to_bigquery import GCSToBigQueryOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryExecuteQueryOperator
from airflow.utils.dates import days_ago

default_args = {
    'owner': 'airflow',
    'start_date': days_ago(1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
    'catchup': False  # 关闭补跑,避免重复同步
}

dag = DAG(
    'archival_dag',
    default_args=default_args,
    schedule_interval='0 0 * * *'
)

# 定义任务变量
BUCKET = 'us-central1-bucket'
FILE_NAME = 'data/{{ds_nodash}}.csv'
TARGET_DATASET = 'testing'
TARGET_TABLE = 'item_bkp_22022023'
TEMP_TABLE = f"{TARGET_DATASET}.temp_item_{{ds_nodash}}"
PROJECT_ID = 'test-dev'

# 1. MySQL到GCS:带数据类型转换的动态查询
def get_sync_query(**context):
    prev_exec_date = context.get('prev_execution_date')
    if not prev_exec_date:
        # 首次运行拉全量,转换数据类型
        return """
            SELECT 
                Id,
                item_type_id,
                user_id,
                TO_BASE64(image_info) AS image_info,  # MySQL BLOB转Base64,适配BigQuery BYTES类型
                text_info,
                location_info,
                item_date_time,
                latitude,
                longitude,
                parcel_id,
                shipment_id,
                transaction_id,
                web_hook_status,
                file_id,
                file_created,
                file_updated,
                pickup_id,
                team_id,
                driver_id,
                container_id
            FROM bkp_item_22_02_2023;
        """
    else:
        # 增量拉取,用调度时间范围过滤
        prev_exec_str = prev_exec_date.strftime('%Y-%m-%d %H:%M:%S')
        return f"""
            SELECT 
                Id,
                item_type_id,
                user_id,
                TO_BASE64(image_info) AS image_info,
                text_info,
                location_info,
                item_date_time,
                latitude,
                longitude,
                parcel_id,
                shipment_id,
                transaction_id,
                web_hook_status,
                file_id,
                file_created,
                file_updated,
                pickup_id,
                team_id,
                driver_id,
                container_id
            FROM bkp_item_22_02_2023 
            WHERE file_updated >= '{prev_exec_str}';
        """

upload_to_gcs = MySQLToGCSOperator(
    task_id='upload_to_gcs',
    sql=get_sync_query,
    bucket=BUCKET,
    filename=FILE_NAME,
    mysql_conn_id='mysql_conn_id',
    gcp_conn_id='google_cloud_default',
    export_format='CSV',
    field_delimiter=',',
    print_header=True,  # 导出带表头,方便BigQuery识别列
    dag=dag
)

# 2. GCS到BigQuery临时表:指定自定义Schema
load_to_temp_bq = GCSToBigQueryOperator(
    task_id='load_to_temp_bq',
    bucket=BUCKET,
    source_objects=[FILE_NAME],
    destination_project_dataset_table=TEMP_TABLE,
    schema_fields=[
        {"name": "Id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "item_type_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "user_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "image_info", "mode": "NULLABLE", "type": "STRING"},
        {"name": "text_info", "mode": "NULLABLE", "type": "STRING"},
        {"name": "location_info", "mode": "NULLABLE", "type": "STRING"},
        {"name": "item_date_time", "mode": "NULLABLE", "type": "TIMESTAMP"},
        {"name": "latitude", "mode": "NULLABLE", "type": "FLOAT"},
        {"name": "longitude", "mode": "NULLABLE", "type": "FLOAT"},
        {"name": "parcel_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "shipment_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "transaction_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "web_hook_status", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "file_id", "mode": "NULLABLE", "type": "STRING"},
        {"name": "file_created", "mode": "NULLABLE", "type": "TIMESTAMP"},
        {"name": "file_updated", "mode": "NULLABLE", "type": "TIMESTAMP"},
        {"name": "pickup_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "team_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "driver_id", "mode": "NULLABLE", "type": "INTEGER"},
        {"name": "container_id", "mode": "NULLABLE", "type": "INTEGER"}
    ],
    write_disposition='WRITE_TRUNCATE',  # 每次覆盖临时表
    create_disposition='CREATE_IF_NEEDED',
    source_format='CSV',
    skip_leading_rows=1,  # 跳过表头
    gcp_conn_id='google_cloud_default',
    dag=dag
)

# 3. 执行MERGE语句实现Upsert
merge_to_target = BigQueryExecuteQueryOperator(
    task_id='merge_to_target',
    sql=f"""
        MERGE `{PROJECT_ID}.{TARGET_DATASET}.{TARGET_TABLE}` AS target
        USING `{PROJECT_ID}.{TEMP_TABLE}` AS source
        ON target.Id = source.Id  # 假设Id是主键,根据实际业务调整
        WHEN MATCHED THEN
            UPDATE SET
                item_type_id = source.item_type_id,
                user_id = source.user_id,
                image_info = FROM_BASE64(source.image_info),  # 转成BigQuery BYTES类型
                text_info = source.text_info,
                location_info = source.location_info,
                item_date_time = source.item_date_time,
                latitude = source.latitude,
                longitude = source.longitude,
                parcel_id = source.parcel_id,
                shipment_id = source.shipment_id,
                transaction_id = source.transaction_id,
                web_hook_status = source.web_hook_status,
                file_id = source.file_id,
                file_created = source.file_created,
                file_updated = source.file_updated,
                pickup_id = source.pickup_id,
                team_id = source.team_id,
                driver_id = source.driver_id,
                container_id = source.container_id
        WHEN NOT MATCHED THEN
            INSERT (
                Id, item_type_id, user_id, image_info, text_info, location_info,
                item_date_time, latitude, longitude, parcel_id, shipment_id,
                transaction_id, web_hook_status, file_id, file_created, file_updated,
                pickup_id, team_id, driver_id, container_id
            )
            VALUES (
                source.Id, source.item_type_id, source.user_id, FROM_BASE64(source.image_info),
                source.text_info, source.location_info, source.item_date_time,
                source.latitude, source.longitude, source.parcel_id, source.shipment_id,
                source.transaction_id, source.web_hook_status, source.file_id,
                source.file_created, source.file_updated, source.pickup_id,
                source.team_id, source.driver_id, source.container_id
            );
        -- 清理临时表
        DROP TABLE IF EXISTS `{PROJECT_ID}.{TEMP_TABLE}`;
    """,
    use_legacy_sql=False,
    gcp_conn_id='google_cloud_default',
    dag=dag
)

# 任务依赖
upload_to_gcs >> load_to_temp_bq >> merge_to_target

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 19:18:32