如何创建订阅使Cloud Function响应BigQuery表新插入条目
解决方案:创建关联Cloud Function的Pub/Sub订阅及BigQuery事件推送配置
要实现BigQuery表staging_tbl插入新数据时触发Cloud Function 2nd Gen并获取新条目,需完成BigQuery事件推送到Pub/Sub主题、创建Pub/Sub订阅关联函数两个核心步骤,具体操作如下:
1. 配置BigQuery向指定Pub/Sub主题推送插入事件
首先要确保BigQuery能将staging_tbl的新插入事件发送到你配置的greeny_data_inserted_in_tbl主题,分两种场景处理:
场景A:流式插入数据(如INSERT语句、流式加载)
- 先给BigQuery服务账号授权:找到项目的BigQuery服务账号(格式为
service-${PROJECT_NUMBER}@gcp-sa-bigquery.iam.gserviceaccount.com),授予其roles/pubsub.publisher权限,目标为你的greeny_data_inserted_in_tbl主题。 - 执行BigQuery SQL配置表的Pub/Sub通知:
ALTER TABLE staging_tbl SET OPTIONS( pubsub_topic = 'projects/${PROJECT_ID}/topics/greeny_data_inserted_in_tbl' );
配置完成后,每次流式插入新行,BigQuery会自动将包含新行数据的消息推送到指定主题。
场景B:批量加载或其他写入方式
通过Cloud Logging路由BigQuery的写入事件到Pub/Sub:
- 进入Cloud Logging控制台,创建新的日志路由,筛选条件设置为:
resource.type="bigquery_dataset" AND protoPayload.methodName="jobs.insert" AND protoPayload.resourceName:"staging_tbl" AND protoPayload.serviceData.jobInsertRequest.job.configuration.query.destinationTable.tableId="staging_tbl" - 将路由目标指定为你的
greeny_data_inserted_in_tbl主题。
2. 创建Pub/Sub订阅关联Cloud Function
你已通过Terraform配置了函数的Pub/Sub触发器,此时GCP会自动为函数生成关联该主题的订阅。如果需要手动创建或用Terraform定义,操作如下:
手动创建(GCP控制台)
- 进入Pub/Sub控制台,找到
greeny_data_inserted_in_tbl主题,点击「创建订阅」。 - 配置参数:
- 订阅ID:自定义(如
greeny_data_inserted_sub) - 交付类型:选择「推送」,推送端点填写你的Cloud Function 2nd Gen URL(格式:
https://${REGION}-${PROJECT_ID}.cloudfunctions.net/${FUNCTION_NAME}) - 重试策略:匹配函数的
RETRY_POLICY_DO_NOT_RETRY,设置为不重试
- 订阅ID:自定义(如
Terraform定义订阅
在现有Terraform配置中添加以下资源:
resource "google_pubsub_subscription" "greeny_data_sub" { name = "greeny_data_inserted_sub" topic = "projects/${var.project_id}/topics/greeny_data_inserted_in_tbl" push_config { push_endpoint = google_cloudfunctions2_function.your_function.service_config.uri } retry_policy { minimum_backoff = "0s" maximum_backoff = "0s" } message_retention_duration = "600s" }
注意:若已通过函数的
event_trigger配置了Pub/Sub主题,无需重复创建订阅,GCP自动生成的订阅已满足需求。
3. 函数内解析新条目示例(Python)
触发函数时,从Pub/Sub消息中提取BigQuery新行数据的示例代码:
import base64 import json def handler(event, context): # 解析Pub/Base64编码的消息 pubsub_message = base64.b64decode(event['data']).decode('utf-8') message_data = json.loads(pubsub_message) # 提取新插入的行数据(适配BigQuery流式通知格式) new_row = message_data.get('rows', [])[0].get('data', {}) column_a = new_row.get('A') column_b = new_row.get('B') column_c = new_row.get('C') # 后续业务逻辑 print(f"新插入条目:A={column_a}, B={column_b}, C={column_c}")
内容的提问来源于stack exchange,提问作者Kara
相关产品推荐
相关产品推荐

