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
相关产品推荐
相关产品推荐

