如何从跨项目Cloud SQL向Airflow DAG中获取数据?
跨项目Cloud SQL数据导入Airflow的可行方案
针对你提到的小数据量(3列6行)场景,完全可以实现将跨项目Cloud SQL的数据导入Airflow,以下是两种直接可行的方案:
方案一:PythonOperator直接查询并获取数据
因为数据量极小,无需中转存储,直接通过PythonOperator连接Cloud SQL(利用共享VPC的私有IP)执行查询,将结果存入Airflow的XCom供后续任务使用,或直接在任务内处理。
步骤:
- 在Composer环境中安装数据库驱动:如果是MySQL类型的Cloud SQL,安装
mysql-connector-python;如果是PostgreSQL,安装psycopg2-binary。 - 编写Python函数完成数据库连接、查询和结果推送。
- 通过PythonOperator调用该函数。
代码示例:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime import mysql.connector def fetch_cloud_sql_data(**context): # 利用共享VPC连接Cloud SQL私有IP conn = mysql.connector.connect( host="<CLOUD_SQL_PRIVATE_IP>", user="<DB_USER>", password="<DB_PASSWORD>", database="<DB_NAME>" ) cursor = conn.cursor(dictionary=True) # 执行查询(按需调整SQL) cursor.execute("SELECT * FROM target_table LIMIT 6") data = cursor.fetchall() # 关闭数据库连接 cursor.close() conn.close() # 将数据推送到XCom,供后续任务调用 context['ti'].xcom_push(key='cloud_sql_data', value=data) print("获取到的Cloud SQL数据:", data) with DAG( dag_id='direct_fetch_cloud_sql', start_date=datetime(2024, 1, 1), schedule_interval='@daily', catchup=False ) as dag: fetch_data = PythonOperator( task_id='fetch_cloud_sql_data', python_callable=fetch_cloud_sql_data, provide_context=True ) fetch_data
方案二:导出到GCS后读取数据
如果偏好先导出到Cloud Storage的方案,可通过CloudSqlExportInstanceOperator完成跨项目导出,再通过GCS Hook读取文件内容到Airflow。
步骤:
- 确保Cloud SQL的服务账号拥有目标项目GCS桶的写入权限。
- 使用
CloudSqlExportInstanceOperator将数据导出为CSV到GCS。 - 通过PythonOperator调用GCS Hook读取并解析CSV数据。
代码示例:
from airflow import DAG from airflow.providers.google.cloud.operators.cloud_sql import CloudSqlExportInstanceOperator from airflow.providers.google.cloud.hooks.gcs import GCSHook from airflow.operators.python import PythonOperator from datetime import datetime import csv from io import StringIO # 配置导出路径(项目A的GCS桶) EXPORT_BUCKET = "your-project-a-bucket" EXPORT_FILE = "sql_export/data.csv" EXPORT_URI = f"gs://{EXPORT_BUCKET}/{EXPORT_FILE}" def read_gcs_exported_data(**context): gcs_hook = GCSHook(gcp_conn_id='google_cloud_default') # 从GCS下载文件内容 file_content = gcs_hook.download(bucket_name=EXPORT_BUCKET, object_name=EXPORT_FILE) # 解析CSV数据 csv_stream = StringIO(file_content.decode('utf-8')) reader = csv.DictReader(csv_stream) data = list(reader) # 推送数据到XCom context['ti'].xcom_push(key='gcs_sql_data', value=data) print("从GCS读取的数据:", data) with DAG( dag_id='cloud_sql_gcs_airflow', start_date=datetime(2024, 1, 1), schedule_interval='@daily', catchup=False ) as dag: export_to_gcs = CloudSqlExportInstanceOperator( task_id='export_cloud_sql_to_gcs', project_id='other-project-id', # Cloud SQL所在项目ID instance='your-cloud-sql-instance', body={ "exportContext": { "fileType": "CSV", "uri": EXPORT_URI, "databases": ["target-db"], "tables": ["target-table"] } }, gcp_conn_id='google_cloud_default' ) read_from_gcs = PythonOperator( task_id='read_gcs_data', python_callable=read_gcs_exported_data, provide_context=True ) export_to_gcs >> read_from_gcs
补充说明:
MySqlOperator仅负责执行SQL并返回执行状态,不会将查询结果带回Airflow,因此不适合你的取数需求。- Airflow虽定位为任务编排,但处理这种轻量级的小数据量取数完全没问题,无需额外的数据处理工具。
内容的提问来源于stack exchange,提问作者Giteshwar Mali
相关产品推荐
相关产品推荐

