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
相关产品推荐
相关产品推荐

