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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 15:58:11