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

Flask API开发中如何在PyCharm读取BigQuery数据到pandas DataFrame并同步至CloudSQL

从BigQuery同步数据到CloudSQL(MySQL)的Flask实现指导

前置准备

  • 确保你的GCP账号已经开通BigQuery和Cloud SQL API权限,并且拥有对应资源的访问角色(BigQuery Data Viewer、BigQuery Job User、Cloud SQL Client)
  • 提前生成GCP服务账号密钥JSON文件,存储到项目可访问路径下
  • 安装所需依赖包:pip install flask google-cloud-bigquery google-cloud-sqlalchemy pymysql sqlalchemy pandas

步骤1:建立BigQuery连接

使用Google官方BigQuery客户端完成连接,只需配置服务账号密钥环境变量即可,无需手动处理连接协议:

import os
import pandas as pd
from google.cloud import bigquery

# 替换为你的服务账号密钥文件实际路径
os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "./gcp-service-account-key.json"
# 初始化BigQuery客户端,完成连接建立
bq_client = bigquery.Client(project="你的GCP项目ID")

可以用测试查询验证连接有效性:

# 测试查询,替换为你的BigQuery表路径
test_query = "SELECT * FROM `项目ID.数据集ID.表名` LIMIT 5"
test_result = bq_client.query(test_query).to_dataframe()
print(test_result)

步骤2:建立CloudSQL(MySQL)连接

推荐用SQLAlchemy统一管理数据库连接,适配Flask生态的同时也能简化后续数据写入逻辑:

from sqlalchemy import create_engine

# 部署到GCP环境(Cloud Run/GKE等)时用Unix套接字连接,无需开放公网IP
# 连接字符串格式:mysql+pymysql://数据库用户名:数据库密码@/数据库名?unix_socket=/cloudsql/项目ID:区域:CloudSQL实例名
cloudsql_conn_str = "mysql+pymysql://my_db_user:my_db_pass@/my_db_name?unix_socket=/cloudsql/gcp-project:us-central1:mysql-instance"

# 本地测试时如果用公网IP连接,替换为以下格式:
# cloudsql_conn_str = "mysql+pymysql://my_db_user:my_db_pass@CloudSQL公网IP:3306/my_db_name"

# 初始化数据库连接引擎
db_engine = create_engine(cloudsql_conn_str)

步骤3:实现数据同步逻辑

优先选择增量同步降低资源消耗,以下是基础同步逻辑示例:

def sync_data():
    # 拉取BigQuery增量数据,此处按更新时间拉取最近1小时的新增/更新数据,可根据业务调整
    sync_query = """
    SELECT id, field1, field2, update_time
    FROM `项目ID.数据集ID.表名`
    WHERE update_time >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR)
    """
    sync_df = bq_client.query(sync_query).to_dataframe()
    
    if sync_df.empty:
        return "无待同步数据"
    
    # 写入CloudSQL,if_exists可根据需求选append/replace,增量同步建议用append+主键冲突更新逻辑
    sync_df.to_sql(
        name="CloudSQL目标表名",
        con=db_engine,
        if_exists="append",
        index=False,
        chunksize=1000, # 数据量较大时分批写入避免超时
        method="multi"
    )
    return f"同步完成,共写入{len(sync_df)}条数据"

步骤4:集成到Flask API接口

from flask import Flask, jsonify

app = Flask(__name__)

@app.route("/api/sync-bq-to-cloudsql", methods=["POST"])
def sync_api():
    result = sync_data()
    return jsonify({"status": "success", "msg": result})

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=8080, debug=False)

生产环境注意事项

  • 敏感信息(服务账号密钥、数据库密码)不要硬编码在代码中,建议用环境变量或密钥管理服务存储
  • 数据量超过1万条时不要在接口同步逻辑中直接处理,容易触发请求超时,建议改为异步任务触发同步
  • 同步前需确保CloudSQL目标表的字段类型、长度和BigQuery拉取的字段匹配,避免写入报错

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 10:15:03