BigQuery DataTransfer Python API凭证适配问题咨询(附代码)
BigQuery数据集复制的Python SDK问题及代码建议
我是GCP项目所有者,在BigQuery UI中可轻松将us多区域的数据集一次性复制为同区域内不同名称的新数据集,但使用Python SDK实现相同操作时遇到以下问题:
- 本地gcloud auth登录无法自动生效,必须设置
GOOGLE_APPLICATION_CREDENTIALS环境变量指向凭证JSON文件。POC场景可行,但代码将部署为Dataflow Pipeline Template起始处的beam.DoFn,不确定该环境下如何配置,期望代码自动使用运行时上下文的服务账号(本地Google auth或GCP中Dataflow服务账号),如同其他BigQuery调用。【更新】使用DataflowRunner时可正常工作并使用预期服务账号,仅开发环境需下载凭证JSON文件。 - 曾出现错误“BigQuery Data Transfer Service does not yet support location: us”,但BigQuery UI中可正常复制。【更新】尝试约15次后问题自行解决,Git历史无相关变更,该问题已不存在。
以下是开发中带注释的代码:
import argparse import logging from datetime import date, timedelta, datetime import time import re import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions from google.cloud import bigquery, bigquery_datatransfer_v1 from google.protobuf.timestamp_pb2 import Timestamp class BackupDataset(beam.DoFn): def __init__(self, bq_analytics_dataset_option, bq_backup_dataset_ttl_days_option, project): self.bq_analytics_dataset = bq_analytics_dataset_option self.bq_backup_dataset_ttl_days = bq_backup_dataset_ttl_days_option self.project = project def create_backup_dataset(self, bq_client, destination_dataset_id): logging.info(f"Creating backup dataset with name: {destination_dataset_id} ...") query_job = bq_client.query(f"CREATE SCHEMA IF NOT EXISTS {destination_dataset_id}") query_job.result() # Wait for job to finish (it usually does immediately, but not always) def create_transfer_config(self, bq_client, bqdt_client, source_dataset_id): date_s = "{:%Y_%m_%d}".format(date.today()) destination_dataset_id = f"{source_dataset_id}_{date_s}" self.create_backup_dataset(bq_client, destination_dataset_id) display_name = f"Backup of {source_dataset_id} on {date_s}" logging.info(f"Creating transfer_config with display name: {display_name} ...") transfer_config = bigquery_datatransfer_v1.TransferConfig( destination_dataset_id=destination_dataset_id, display_name=display_name, data_source_id="cross_region_copy", params={ "source_project_id": self.project, "source_dataset_id": source_dataset_id, "overwrite_destination_table": True }, schedule_options={ "disable_auto_scheduling": True } # run this only once - the default is recurring ) # In order for this call to not puke with "Failed to find a valid credential. The field 'version_info' or 'service_account_name' must be specified.", # the GOOGLE_APPLICATION_CREDENTIALS ENV var must be set to point to a downloaded service account key file on the local machine. # Despite the error text, this call does NOT accept a `service_account_name` param. # FIXME: Find a way for this to work without ^^ for ease of developer maintenance remote_transfer_config = bqdt_client.create_transfer_config( parent=bqdt_client.common_project_path(project), transfer_config=transfer_config ) logging.info(f"Created transfer_config with name: {remote_transfer_config.name} ...") return remote_transfer_config def run_transfer_config(self, bqdt_client, config): logging.info(f"Running transfer config with name {config.name} ...") start_time = Timestamp(seconds=int(time.time())) request = bigquery_datatransfer_v1.types.StartManualTransferRunsRequest( { "parent": config.name, "requested_run_time": start_time } ) bqdt_client.start_manual_transfer_runs(request, timeout=360) # ...
代码优化建议
1. 开发环境凭证自动适配
本地开发时无需强制设置环境变量,可通过显式加载gcloud默认凭证解决:
from google.auth import default # 在初始化客户端时添加凭证加载逻辑 credentials, project_id = default() bqdt_client = bigquery_datatransfer_v1.DataTransferServiceClient(credentials=credentials)
这样本地gcloud auth登录的凭证会被自动加载,无需手动指定JSON文件路径。
2. 数据集创建优化
替换SQL创建方式为BigQuery SDK原生方法,更易扩展参数(如TTL、位置):
def create_backup_dataset(self, bq_client, destination_dataset_id, source_dataset_id): logging.info(f"Creating backup dataset with name: {destination_dataset_id} ...") project_id, dataset_id = destination_dataset_id.split('.', 1) dataset_ref = bigquery.DatasetReference(project_id, dataset_id) dataset = bigquery.Dataset(dataset_ref) # 继承源数据集位置,避免区域不匹配 source_dataset = bq_client.get_dataset(source_dataset_id) dataset.location = source_dataset.location # 设置默认表过期时间 if self.bq_backup_dataset_ttl_days: dataset.default_table_expiration_ms = self.bq_backup_dataset_ttl_days * 86400 * 1000 bq_client.create_dataset(dataset, exists_ok=True)
调用时传入源数据集ID,确保备份数据集位置与源一致。
3. 传输配置显式指定位置
创建传输配置时,使用带区域的父路径,避免潜在的区域适配问题:
# 替换原parent参数 parent = bqdt_client.common_location_path(self.project, "us") remote_transfer_config = bqdt_client.create_transfer_config( parent=parent, transfer_config=transfer_config )
BigQuery Data Transfer Service资源为区域化,显式指定位置可避免默认值不匹配。
4. 传输任务状态检查
添加任务状态轮询,确保备份任务成功完成:
def run_transfer_config(self, bqdt_client, config): logging.info(f"Running transfer config with name {config.name} ...") start_time = Timestamp(seconds=int(time.time())) request = bigquery_datatransfer_v1.types.StartManualTransferRunsRequest( {"parent": config.name, "requested_run_time": start_time} ) runs = bqdt_client.start_manual_transfer_runs(request, timeout=360) # 轮询任务状态 for run in runs.runs: while True: run_status = bqdt_client.get_transfer_run(run.name) if run_status.state in [bigquery_datatransfer_v1.TransferState.SUCCEEDED, bigquery_datatransfer_v1.TransferState.FAILED]: break time.sleep(10) if run_status.state == bigquery_datatransfer_v1.TransferState.FAILED: logging.error(f"Transfer failed: {run_status.error_status.message}") raise RuntimeError(f"Dataset backup failed: {run_status.error_status.message}") logging.info(f"Transfer run {run.name} completed successfully")
内容的提问来源于stack exchange,提问作者aec
相关产品推荐
相关产品推荐

