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

如何使用Airflow将BigQuery表迁移至PostgreSQL表

BigQuery到PostgreSQL的Airflow迁移方案及配置指南

可行的迁移方案

方案1:BigQuery导出GCS + Python自定义导入PostgreSQL

这是最轻量化的方案,无需额外依赖其他GCP服务,利用现有Airflow组件即可完成:

  • 用BigQueryToGCSOperator将BigQuery表导出为CSV/Parquet/Avro格式到GCS。
  • 用PythonOperator编写自定义逻辑,通过GCS SDK下载文件,再用psycopg2或pandas写入PostgreSQL。适合中小数据量迁移。

方案2:Dataflow端到端迁移

利用Dataflow的Apache Beam框架直接读取BigQuery数据并写入PostgreSQL,适合大数据量的流式或批量迁移:

  • 编写Apache Beam作业,用ReadFromBigQuery读取数据,WriteToJdbc写入PostgreSQL。
  • 通过Airflow的DataflowSubmitJobOperator提交作业到GCP。你的Provisioner角色可以创建并配置Dataflow所需的服务账号权限。

方案3:DataProc Spark作业迁移

借助DataProc的Spark集群处理大数据量迁移,Spark对BQ和PostgreSQL的兼容性较好:

  • 编写PySpark作业,用Spark的BQ连接器读取数据,通过JDBC驱动写入PostgreSQL。
  • 通过Airflow的Dataproc系列Operator创建集群、提交作业、完成后销毁集群,节省成本。

Airflow配置步骤

1. 安装必要依赖包

在Airflow部署环境中执行以下命令安装提供者包:

pip install apache-airflow-providers-google apache-airflow-providers-postgres

2. 配置GCP连接

  • 进入Airflow UI的「Admin > Connections」页面,新建连接:
    • 连接类型选择「Google Cloud」
    • 认证方式选择「Keyfile JSON」,上传拥有GCS、BigQuery操作权限的服务账号密钥
    • 填写GCP项目ID,测试连接确认可用

3. 配置PostgreSQL连接

  • 同样在「Admin > Connections」页面新建连接:
    • 连接类型选择「Postgres」
    • 填写PostgreSQL的主机地址、端口、数据库名、用户名、密码
    • 点击「Test Connection」确保连接正常

方案示例代码

方案1:GCS中转 + Python导入的DAG

from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryToGCSOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.operators.python import PythonOperator
from google.cloud import storage
import pandas as pd
from datetime import datetime

def load_gcs_to_postgres():
    # 初始化GCS客户端并下载文件
    storage_client = storage.Client()
    bucket = storage_client.bucket("your-gcs-bucket-name")
    blob = bucket.blob("bq-exports/your-table.csv")
    temp_file = "/tmp/bq_export.csv"
    blob.download_to_filename(temp_file)

    # 写入PostgreSQL(大数据量建议用COPY命令提升效率)
    pg_hook = PostgresHook(postgres_conn_id="your-postgres-connection-id")
    conn = pg_hook.get_conn()
    df = pd.read_csv(temp_file)
    df.to_sql("target_postgres_table", conn, if_exists="replace", index=False)
    
    conn.commit()
    conn.close()

default_args = {'start_date': datetime(2024, 1, 1)}

with DAG('bq_to_postgres_via_gcs', default_args=default_args, schedule_interval=None) as dag:
    export_bq = BigQueryToGCSOperator(
        task_id='export_bq_to_gcs',
        source_project_dataset_table='your-gcp-project.your-dataset.your-table',
        destination_cloud_storage_uris=['gs://your-gcs-bucket-name/bq-exports/your-table.csv'],
        export_format='CSV',
        print_header=True,
        gcp_conn_id="your-gcp-connection-id"
    )

    load_pg = PythonOperator(
        task_id='load_gcs_to_postgres',
        python_callable=load_gcs_to_postgres
    )

    export_bq >> load_pg

方案2:Dataflow端到端迁移的DAG

先编写Apache Beam作业文件(上传至GCS):

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions
from apache_beam.io.gcp.bigquery import ReadFromBigQuery
from apache_beam.io.jdbc import WriteToJdbc

class BQToPGOptions(PipelineOptions):
    @classmethod
    def _add_argparse_args(cls, parser):
        parser.add_argument('--bq_table', help='BigQuery表路径,格式:项目.数据集.表名')
        parser.add_argument('--pg_jdbc_url', help='PostgreSQL JDBC URL,格式:jdbc:postgresql://主机:端口/数据库')
        parser.add_argument('--pg_table', help='PostgreSQL目标表名')
        parser.add_argument('--pg_user', help='PostgreSQL用户名')
        parser.add_argument('--pg_password', help='PostgreSQL密码')

def run():
    options = BQToPGOptions()
    gcp_options = options.view_as(GoogleCloudOptions)
    gcp_options.project = 'your-gcp-project'
    gcp_options.region = 'us-central1'
    gcp_options.job_name = 'bq-to-pg-migration'
    gcp_options.staging_location = 'gs://your-gcs-bucket/staging'
    gcp_options.temp_location = 'gs://your-gcs-bucket/temp'

    with beam.Pipeline(options=options) as p:
        (p
         | '读取BQ数据' >> ReadFromBigQuery(table=options.bq_table)
         | '转换为Row格式' >> beam.Map(lambda x: beam.Row(**x))
         | '写入PostgreSQL' >> WriteToJdbc(
             table_name=options.pg_table,
             jdbc_url=options.pg_jdbc_url,
             driver_class_name='org.postgresql.Driver',
             username=options.pg_user,
             password=options.pg_password
         )
        )

if __name__ == '__main__':
    run()

Airflow DAG提交Dataflow作业:

from airflow.providers.google.cloud.operators.dataflow import DataflowSubmitJobOperator
from datetime import datetime

default_args = {'start_date': datetime(2024, 1, 1)}

with DAG('bq_to_postgres_dataflow', default_args=default_args, schedule_interval=None) as dag:
    submit_dataflow = DataflowSubmitJobOperator(
        task_id='submit_dataflow_job',
        job_name='bq-to-pg-migration',
        project_id='your-gcp-project',
        region='us-central1',
        job_type='PYTHON',
        python_file='gs://your-gcs-bucket/beam-jobs/bq_to_pg.py',
        runtime_environment={'tempLocation': 'gs://your-gcs-bucket/temp'},
        parameters={
            'bq_table': 'your-gcp-project.your-dataset.your-table',
            'pg_jdbc_url': 'jdbc:postgresql://your-pg-host:5432/your-db',
            'pg_table': 'target_table',
            'pg_user': 'your-pg-user',
            'pg_password': 'your-pg-password'
        },
        gcp_conn_id='your-gcp-connection-id'
    )

方案3:DataProc Spark迁移的DAG

先编写PySpark作业文件(上传至GCS):

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("BQToPostgres").getOrCreate()

# 读取BigQuery数据
df = spark.read.format("bigquery") \
    .option("table", "your-gcp-project.your-dataset.your-table") \
    .load()

# 写入PostgreSQL(需提前将PostgreSQL JDBC驱动上传至GCS)
df.write.format("jdbc") \
    .option("url", "jdbc:postgresql://your-pg-host:5432/your-db") \
    .option("dbtable", "target_table") \
    .option("user", "your-pg-user") \
    .option("password", "your-pg-password") \
    .option("driver", "org.postgresql.Driver") \
    .mode("overwrite") \
    .save()

spark.stop()

Airflow DAG管理DataProc集群和作业:

from airflow.providers.google.cloud.operators.dataproc import DataprocCreateClusterOperator, DataprocSubmitPySparkJobOperator, DataprocDeleteClusterOperator
from datetime import datetime

CLUSTER_NAME = 'bq-to-pg-cluster'
PROJECT_ID = 'your-gcp-project'
REGION = 'us-central1'

default_args = {'start_date': datetime(2024, 1, 1)}

with DAG('bq_to_postgres_dataproc', default_args=default_args, schedule_interval=None) as dag:
    create_cluster = DataprocCreateClusterOperator(
        task_id='create_dataproc_cluster',
        project_id=PROJECT_ID,
        region=REGION,
        cluster_name=CLUSTER_NAME,
        cluster_config={
            'master_config': {'num_instances': 1, 'machine_type_uri': 'n1-standard-2'},
            'worker_config': {'num_instances': 2, 'machine_type_uri': 'n1-standard-2'}
        },
        gcp_conn_id='your-gcp-connection-id'
    )

    submit_spark_job = DataprocSubmitPySparkJobOperator(
        task_id='submit_spark_job',
        project_id=PROJECT_ID,
        region=REGION,
        cluster_name=CLUSTER_NAME,
        main='gs://your-gcs-bucket/spark-jobs/bq_to_pg.py',
        jars=['gs://your-gcs-bucket/jars/postgresql-42.7.3.jar']
    )

    delete_cluster = DataprocDeleteClusterOperator(
        task_id='delete_dataproc_cluster',
        project_id=PROJECT_ID,
        region=REGION,
        cluster_name=CLUSTER_NAME,
        trigger_rule='all_done',
        gcp_conn_id='your-gcp-connection-id'
    )

    create_cluster >> submit_spark_job >> delete_cluster

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 03:35:03