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

如何从跨项目Cloud SQL向Airflow DAG中获取数据?

跨项目Cloud SQL数据导入Airflow的可行方案

针对你提到的小数据量(3列6行)场景,完全可以实现将跨项目Cloud SQL的数据导入Airflow,以下是两种直接可行的方案:

方案一:PythonOperator直接查询并获取数据

因为数据量极小,无需中转存储,直接通过PythonOperator连接Cloud SQL(利用共享VPC的私有IP)执行查询,将结果存入Airflow的XCom供后续任务使用,或直接在任务内处理。

步骤:

  1. 在Composer环境中安装数据库驱动:如果是MySQL类型的Cloud SQL,安装mysql-connector-python;如果是PostgreSQL,安装psycopg2-binary。
  2. 编写Python函数完成数据库连接、查询和结果推送。
  3. 通过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。

步骤:

  1. 确保Cloud SQL的服务账号拥有目标项目GCS桶的写入权限。
  2. 使用CloudSqlExportInstanceOperator将数据导出为CSV到GCS。
  3. 通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 05:35:06