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

如何在Cloud Composer的Airflow任务中基于BigQuery表动态创建分国家表?

实现思路

核心逻辑分为三步:

  1. 获取增量数据范围与国家列表:记录上次同步时间,从源表中提取本次需要处理的国家(仅包含有新增数据的国家)
  2. 动态生成任务:利用Airflow的动态任务映射(Dynamic Task Mapping)为每个国家创建独立处理任务
  3. 建表/追加数据:对每个国家,判断目标表是否存在,不存在则创建并插入增量数据;存在则直接追加新增数据
具体实现代码

假设源表为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
关键细节说明
  1. 增量同步逻辑:
    • 依赖源表的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`)
      
  2. 动态任务映射:
    • 基于Airflow 2.2+版本的expand特性,运行时根据国家列表动态生成任务,无需提前定义固定任务
    • 旧版本Airflow可改用循环生成PythonOperator或BigQueryInsertJobOperator实例
  3. 权限配置:
    • Cloud Composer的服务账号需具备BigQuery Data Editor角色,或细粒度的bigquery.tables.create、bigquery.tables.updateData权限
  4. 表命名规范:
    • 给国家表加country_前缀,避免国家代码(如US、UK)与SQL关键字冲突
特殊场景处理
  • 已移除的国家:本逻辑仅处理当前源表中存在的国家,对已从源表移除的国家不做任何操作,符合需求
  • 空数据场景:如果某个国家本次没有新增数据,fetch_target_countries不会将其纳入列表,对应任务不会执行,避免无效操作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 02:35:16