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

无需Cloud Functions,实现GCP Cloud Composer中Airflow DAG自动触发

替代Cloud Functions触发Airflow DAG的方案

方案1:Cloud Run + Pub/Sub + Airflow REST API

  • 配置GCS存储桶的对象创建事件,将事件推送到指定的Pub/Sub主题
  • 部署一个Cloud Run服务,订阅该Pub/Sub主题,服务逻辑为:收到GCS对象创建事件后,检查是否是.csv文件,再调用Airflow的REST API触发目标DAG
  • 关键配置:
    • 给Cloud Run服务分配IAM权限,允许它调用Airflow API(需Airflow的DAG触发权限)
    • 给Pub/Sub主题分配权限,允许GCS推送事件、Cloud Run订阅主题
  • 示例代码片段(Python):
import os
import requests
import base64
import json
from flask import Flask, request

app = Flask(__name__)

AIRFLOW_API_URL = os.environ.get("AIRFLOW_API_URL")
AIRFLOW_AUTH_TOKEN = os.environ.get("AIRFLOW_AUTH_TOKEN")

@app.route('/', methods=['POST'])
def trigger_dag():
    envelope = request.get_json()
    if not envelope or 'message' not in envelope:
        return ('Bad Request', 400)
    
    message = envelope['message']
    data = message.get('data')
    if not data:
        return ('Bad Request', 400)
    
    event_data = base64.b64decode(data).decode('utf-8')
    event = json.loads(event_data)
    if event['name'].endswith('.csv'):
        headers = {
            'Authorization': f'Bearer {AIRFLOW_AUTH_TOKEN}',
            'Content-Type': 'application/json'
        }
        payload = {
            "conf": {"gcs_file_path": f"gs://{event['bucket']}/{event['name']}"}
        }
        requests.post(
            f"{AIRFLOW_API_URL}/api/v1/dags/your_dag_id/dagRuns",
            headers=headers,
            json=payload
        )
        return ('OK', 200)
    return ('Ignored', 200)

if __name__ == '__main__':
    app.run(host='0.0.0.0', port=int(os.environ.get('PORT', 8080)))

方案2:Airflow内置GCSSensor轮询触发

  • 在Airflow中创建带传感器的DAG,用GCSSensor监听指定GCS桶中的.csv文件
  • 配置传感器的poke_interval(轮询间隔),比如每5分钟检查一次,平衡实时性和资源消耗
  • 传感器发现新CSV后触发后续任务,同时将已处理文件移至归档目录避免重复触发
  • 示例DAG代码片段:
from airflow import DAG
from airflow.providers.google.cloud.sensors.gcs import GCSSensor
from airflow.providers.google.cloud.operators.gcs import GCSMoveObjectOperator
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'gcs_csv_trigger_dag',
    default_args=default_args,
    description='Trigger on new CSV in GCS',
    schedule_interval=timedelta(minutes=5),
    catchup=False,
) as dag:

    wait_for_csv = GCSSensor(
        task_id='wait_for_new_csv',
        bucket_name='your-bucket',
        object_key='*.csv',
        wildcard_match=True,
        poke_interval=300,
        mode='reschedule'
    )

    def process_csv(**context):
        file_path = context['ti'].xcom_pull(task_ids='wait_for_new_csv')
        # 写入CSV处理逻辑
        print(f"Processing file: {file_path}")

    process_task = PythonOperator(
        task_id='process_csv_file',
        python_callable=process_csv,
        provide_context=True
    )

    archive_file = GCSMoveObjectOperator(
        task_id='archive_processed_csv',
        source_bucket='your-bucket',
        source_object="{{ ti.xcom_pull(task_ids='wait_for_new_csv') }}",
        destination_bucket='your-bucket',
        destination_object="processed/{{ ti.xcom_pull(task_ids='wait_for_new_csv').split('/')[-1] }}"
    )

    wait_for_csv >> process_task >> archive_file

方案3:Cloud Scheduler + Airflow REST API

  • 用Cloud Scheduler创建定时任务,定期检查GCS桶中是否有未处理的.csv文件
  • 定时任务直接调用Airflow的REST API触发DAG,或先执行gcloud storage ls筛选新文件再触发
  • 配置要点:
    • 给Cloud Scheduler分配IAM权限,允许它调用Airflow API和执行GCS命令
    • 根据业务需求设置调度频率,比如每10分钟一次
  • 示例Cloud Scheduler配置(HTTP目标):
    • 目标URL:https://your-composer-webserver-url/api/v1/dags/your_dag_id/dagRuns
    • 请求方法:POST
    • 请求体:{"conf": {"check_new_csv": true}}
    • 授权方式:OAuth 2.0,选择具备Airflow触发权限的服务账号

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 06:23:12