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

Pub/Sub接收不定长JSON数组 如何通过Dataflow SQL导出至BigQuery

Pub/Sub Schema适配JSON数组的定义方法

你可以直接在Pub/Sub中创建JSON类型的Schema,顶层定义为数组类型,items字段指定数组内对象的统一结构即可,示例Schema如下:

{
  "type": "array",
  "items": {
    "type": "object",
    "properties": {
      "name": { "type": "string" },
      "age": { "type": "integer" },
      "car": { "type": "string" }
    },
    "required": ["name", "age", "car"]
  }
}

将该Schema绑定到你的Pub/Sub主题后,所有发布到主题的消息都会自动校验格式,只有符合数组结构的消息才能成功写入。

Dataflow SQL遍历数组导出到BigQuery的实现

完全可以通过Dataflow SQL处理不定长JSON数组,核心是使用UNNEST函数将数组展开为多行,再写入BigQuery,两种常见实现方式如下:

方式一:基于已绑定Pub/Sub Schema的查询

如果你的主题已经绑定了上述Schema,在Dataflow SQL中定义源表时,数组字段会自动识别为ARRAY<STRUCT>类型,直接展开即可:

SELECT
  -- 读取展开后的单条用户数据
  user_item.name,
  user_item.age,
  user_item.car,
  -- 可选保留Pub/Sub元数据,如消息发布时间
  event_timestamp AS msg_publish_time
FROM
  -- 替换为你在Dataflow SQL中映射的Pub/Sub源表名
  `your_pubsub_source_table`,
  UNNEST(user_list) AS user_item

执行该查询后将结果Sink到目标BigQuery表即可,不管数组长度是多少,都会自动拆分为独立行写入。

方式二:不绑定Schema,直接解析消息字符串

如果你不想预先在Pub/Sub绑定Schema,也可以将消息载荷当做字符串读取,用JSON解析函数手动解析后展开:

SELECT
  JSON_VALUE(user_item, '$.name') AS name,
  CAST(JSON_VALUE(user_item, '$.age') AS INT64) AS age,
  JSON_VALUE(user_item, '$.car') AS car,
  event_timestamp AS msg_publish_time
FROM (
  SELECT
    -- 将Pub/Sub二进制载荷转为字符串
    CAST(payload AS STRING) AS msg_str,
    event_timestamp
  FROM `your_pubsub_source_table`
),
-- 解析字符串为JSON数组并展开
UNNEST(JSON_EXTRACT_ARRAY(msg_str)) AS user_item

该方式灵活度更高,适合数组内对象结构可能频繁调整的场景。

补充说明
  • 提前绑定Pub/Sub Schema的方案自带消息校验能力,可以提前拦截非法格式的消息,避免下游SQL执行报错
  • 直接解析字符串的方案无需修改Pub/Sub配置,结构调整时仅需修改SQL逻辑即可,迭代成本更低

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 07:45:06