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
相关产品推荐
相关产品推荐

