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

求助:将Customer.io Webhook数据写入BigQuery的Google Cloud Function Python代码

解决方案:将Customer.io Webhook数据写入BigQuery

前置准备

  • 确保你的Google Cloud Function已启用BigQuery API,且云函数关联的服务账号拥有BigQuery Data Editor或BigQuery Editor权限(至少具备目标表的写入权限)。
  • 提前在BigQuery中创建匹配数据结构的目标表,参考以下结构定义:

BigQuery表结构SQL语句

CREATE TABLE `你的项目ID.你的数据集ID.你的表名` (
  event_id STRING,
  metric STRING,
  object_type STRING,
  timestamp TIMESTAMP,
  action_id INTEGER,
  campaign_id INTEGER,
  customer_id STRING,
  delivery_id STRING,
  recipient STRING,
  subject STRING,
  identifier_id STRING
);

云函数配置与代码修改

1. 依赖配置

在云函数的requirements.txt文件中添加BigQuery客户端库:

google-cloud-bigquery==3.11.4

(版本可选用最新稳定版)

2. 修改后的Python代码

import os
from google.cloud import bigquery
from datetime import datetime

def webhook(request):
    # 初始化BigQuery客户端
    client = bigquery.Client()
    # 替换为你的项目、数据集、表ID
    table_id = "你的项目ID.你的数据集ID.你的表名"

    request_json = request.get_json()
    if request.method == 'POST' and request_json:
        try:
            # 解析Webhook数据字段
            data = request_json.get('data', {})
            # 将Unix时间戳转换为BigQuery兼容的TIMESTAMP格式
            event_timestamp = datetime.fromtimestamp(request_json['timestamp']).isoformat()

            # 构造插入数据行
            row_to_insert = {
                "event_id": request_json['event_id'],
                "metric": request_json['metric'],
                "object_type": request_json['object_type'],
                "timestamp": event_timestamp,
                "action_id": data.get('action_id'),
                "campaign_id": data.get('campaign_id'),
                "customer_id": data.get('customer_id'),
                "delivery_id": data.get('delivery_id'),
                "recipient": data.get('recipient'),
                "subject": data.get('subject'),
                "identifier_id": data.get('identifiers', {}).get('id')
            }

            # 写入BigQuery
            errors = client.insert_rows_json(table_id, [row_to_insert])
            if not errors:
                print("数据成功写入BigQuery")
                return 'success'
            else:
                print(f"写入错误: {errors}")
                return 'failed'
        except Exception as e:
            print(f"请求处理出错: {str(e)}")
            return 'failed'
    else:
        return 'failed'

关键说明

  • 字段容错:使用get()方法提取字段,避免因Webhook数据缺失部分字段导致函数崩溃。
  • 时间格式转换:原始数据中的Unix时间戳需转换为ISO格式字符串,BigQuery会自动解析为TIMESTAMP类型。
  • 权限排查:若出现写入权限报错,检查云函数服务账号的IAM角色配置,确保拥有BigQuery写入权限。

内容的提问来源于stack exchange,提问作者Zeus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 23:31:12