无需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触发权限的服务账号
- 目标URL:
内容的提问来源于stack exchange,提问作者user23605998
相关产品推荐
相关产品推荐

