如何在Cloud Composer的Airflow任务中基于BigQuery表动态创建分国家表?
实现思路
核心逻辑分为三步:
- 获取增量数据范围与国家列表:记录上次同步时间,从源表中提取本次需要处理的国家(仅包含有新增数据的国家)
- 动态生成任务:利用Airflow的动态任务映射(Dynamic Task Mapping)为每个国家创建独立处理任务
- 建表/追加数据:对每个国家,判断目标表是否存在,不存在则创建并插入增量数据;存在则直接追加新增数据
具体实现代码
假设源表为your-project.your-dataset.source_table,目标表命名为your-project.your-dataset.country_{country_code}(加前缀避免SQL关键字冲突),以下是完整的Airflow DAG示例:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator from airflow.models import Variable from airflow.utils.dates import days_ago from datetime import timedelta from google.cloud import bigquery # 默认参数配置 default_args = { 'owner': 'airflow', 'depends_on_past': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } def fetch_target_countries(**context): """获取本次需要处理的国家列表,并更新同步时间戳""" client = bigquery.Client() # 从Airflow变量获取上次同步时间,初始值设为1970-01-01 last_sync = Variable.get('bq_country_sync_timestamp', default_var='1970-01-01 00:00:00') # 查询有新增数据的国家(假设源表有update_time字段记录数据更新时间) query = f""" SELECT DISTINCT country_code FROM `your-project.your-dataset.source_table` WHERE update_time > TIMESTAMP('{last_sync}') """ query_job = client.query(query) country_list = [row.country_code for row in query_job.result()] # 更新同步时间为当前DAG执行时间 current_sync = context['execution_date'].isoformat() Variable.set('bq_country_sync_timestamp', current_sync) # 将国家列表推送到XCom供后续任务使用 context['ti'].xcom_push(key='country_list', value=country_list) return country_list with DAG( 'dynamic_bq_country_table_sync', default_args=default_args, description='动态创建/追加BigQuery国家维度表', schedule_interval='@daily', # 根据业务需求调整调度频率 start_date=days_ago(1), catchup=False, tags=['bigquery', 'cloud-composer'] ) as dag: # 任务1:获取需要处理的国家列表 get_countries = PythonOperator( task_id='fetch_target_countries', python_callable=fetch_target_countries, provide_context=True, execution_timeout=timedelta(minutes=10) ) # 任务2:动态为每个国家执行建表/追加操作 sync_country_table = BigQueryInsertJobOperator.partial( task_id='sync_country_table', gcp_conn_id='google_cloud_default', # 确保已配置GCP连接 configuration={ "query": { "query": """ -- 创建目标表(如果不存在) CREATE TABLE IF NOT EXISTS `your-project.your-dataset.country_{{ params.country_code }}` ( value INT64, inserted_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP() ); -- 插入增量数据 INSERT INTO `your-project.your-dataset.country_{{ params.country_code }}` (value) SELECT value FROM `your-project.your-dataset.source_table` WHERE country_code = '{{ params.country_code }}' AND update_time > TIMESTAMP('{{ params.last_sync }}'); """, "useLegacySql": False } } ).expand( params=[ { "country_code": country, "last_sync": Variable.get('bq_country_sync_timestamp') } for country in "{{ ti.xcom_pull(key='country_list', task_ids='fetch_target_countries') }}" ] ) # 设置任务依赖 get_countries >> sync_country_table
关键细节说明
- 增量同步逻辑:
- 依赖源表的
update_time字段识别新增数据,如果源表没有该字段,可改用主键去重:INSERT INTO `target_table` (value, id) SELECT value, id FROM `source_table` WHERE country_code = 'XX' AND id NOT IN (SELECT id FROM `target_table`)
- 依赖源表的
- 动态任务映射:
- 基于Airflow 2.2+版本的
expand特性,运行时根据国家列表动态生成任务,无需提前定义固定任务 - 旧版本Airflow可改用循环生成
PythonOperator或BigQueryInsertJobOperator实例
- 基于Airflow 2.2+版本的
- 权限配置:
- Cloud Composer的服务账号需具备
BigQuery Data Editor角色,或细粒度的bigquery.tables.create、bigquery.tables.updateData权限
- Cloud Composer的服务账号需具备
- 表命名规范:
- 给国家表加
country_前缀,避免国家代码(如US、UK)与SQL关键字冲突
- 给国家表加
特殊场景处理
- 已移除的国家:本逻辑仅处理当前源表中存在的国家,对已从源表移除的国家不做任何操作,符合需求
- 空数据场景:如果某个国家本次没有新增数据,
fetch_target_countries不会将其纳入列表,对应任务不会执行,避免无效操作
内容的提问来源于stack exchange,提问作者drake10k
相关产品推荐
相关产品推荐

