如何在DataTransferServiceClient中集成TransferConfig并触发S3至BigQuery手动传输
手动触发S3到BigQuery数据传输:TransferConfig集成与config_id获取
1. 创建TransferConfig(若未配置)
如果还没有现成的S3到BigQuery传输配置,你需要先创建TransferConfig对象,明确数据源、目标及传输规则:
- 核心参数说明:
data_source_id:固定为s3(S3数据源标识)display_name:自定义配置名称,方便后续识别params:包含S3路径、AWS凭证(建议用Secret Manager存储)、文件格式等关键信息destination_dataset_id:目标BigQuery数据集IDschedule:无需定时触发时设为"none"
示例代码:
from google.cloud import bigquery_datatransfer_v1 from google.protobuf.struct_pb2 import Struct import time client = bigquery_datatransfer_v1.DataTransferServiceClient() # 构建父资源路径,格式为projects/{project_id}/locations/{location} parent = client.common_project_path("your-project-id", "us-central1") # 配置S3传输参数 transfer_params = Struct() transfer_params["s3_source_path"] = "s3://your-bucket/path/to/data/" transfer_params["aws_access_key_id"] = "your-aws-access-key" transfer_params["aws_secret_access_key"] = "your-aws-secret-key" transfer_params["file_format"] = "CSV" transfer_params["skip_leading_rows"] = 1 # 若CSV含表头需设置 # 定义TransferConfig transfer_config = bigquery_datatransfer_v1.TransferConfig( data_source_id="s3", display_name="S3 to BQ Manual Transfer", params=transfer_params, destination_dataset_id="your-bq-dataset", schedule="none" ) # 发送创建请求 response = client.create_transfer_config(parent=parent, transfer_config=transfer_config) # 提取config_id(从返回的name字段中分割获取) config_id = response.name.split("/")[-1] print(f"Created Transfer Config ID: {config_id}")
2. 获取已有TransferConfig的config_id
如果已存在传输配置,可通过list_transfer_configs方法查询并筛选目标配置:
client = bigquery_datatransfer_v1.DataTransferServiceClient() parent = client.common_project_path("your-project-id", "us-central1") # 遍历所有传输配置,筛选S3类型的目标配置 for config in client.list_transfer_configs(parent=parent): if config.data_source_id == "s3" and config.display_name == "S3 to BQ Manual Transfer": config_id = config.name.split("/")[-1] print(f"Found Transfer Config ID: {config_id}") break
config.name的格式为projects/{project}/locations/{location}/transferConfigs/{config_id},通过分割字符串即可提取末尾的config_id。
3. 手动触发传输
拿到config_id后,调用start_manual_transfer_runs方法触发即时传输:
client = bigquery_datatransfer_v1.DataTransferServiceClient() # 构建TransferConfig的完整路径 transfer_config_name = client.transfer_config_path( "your-project-id", "us-central1", config_id ) # 触发手动传输 response = client.start_manual_transfer_runs( parent=transfer_config_name, requested_run_time={"seconds": int(time.time())} # 指定当前时间为触发时间 ) print(f"Triggered manual transfer run: {response.run_name}")
关键注意事项
- 确保服务账号拥有
bigquery.datatransfers.transferConfigs.create和bigquery.datatransfers.startManualTransferRuns权限 - AWS凭证建议通过Google Cloud Secret Manager存储,避免硬编码(可在
params中指定aws_secret_access_key_secret_resource为Secret Manager资源路径) - 传输参数需与S3文件格式严格匹配(如CSV需指定分隔符、是否跳过表头)
内容的提问来源于stack exchange,提问作者Joon
相关产品推荐
相关产品推荐

