基于BigQuery + dbt动态解析JSON并生成目标表
解决方案:用dbt+BigQuery动态创建包含变化字段的表
问题背景
给定JSON消息结构如下:
{"data": {"schema":"dev","payload": {"lastmodifieddate": "2022-11-12 00:01:28","changeeventheader": {"changetype": "UPDATE","changefields": ["lastmodifieddate","product_value"],"committimestamp": 18478596845860,"recordIds":["568069"]},"product_value" : 20000}}}
需要在BigQuery中通过dbt动态生成表,表结构包含recordIds、changetype以及changefields中指定的动态字段,最终输出表结构如下:
| recordIds | product_value | lastmodifieddate | changetype |
|---|---|---|---|
| 568069 | 20000 | 2022-11-12 00:01:28 | UPDATE |
实现步骤
1. 核心思路
通过BigQuery原生JSON解析函数提取嵌套数据,结合dbt的Jinja模板能力动态识别changefields中的字段列表,避免硬编码字段名,实现表结构的动态适配。
2. dbt模型代码实现
假设源数据存储在BigQuery表{{ source('raw', 'change_events') }}中,其中json_payload列存储完整的JSON字符串:
{% set changefields_query %} SELECT DISTINCT TRIM(JSON_EXTRACT_SCALAR(field, '$')) AS field FROM {{ source('raw', 'change_events') }}, UNNEST(JSON_EXTRACT_ARRAY(json_payload, '$.data.payload.changeeventheader.changefields')) AS field {% endset %} {% set results = run_query(changefields_query) %} {% if execute %} {% set changefields = results.columns[0].values() %} {% else %} {% set changefields = [] %} {% endif %} WITH parsed_data AS ( SELECT -- 将recordIds数组拆分为单行记录 UNNEST(JSON_EXTRACT_ARRAY(json_payload, '$.data.payload.changeeventheader.recordIds')) AS recordIds, -- 提取变更类型 JSON_EXTRACT_SCALAR(json_payload, '$.data.payload.changeeventheader.changetype') AS changetype, -- 动态生成变化字段的解析逻辑 {% for field in changefields %} {% if field == 'lastmodifieddate' %} PARSE_DATETIME('%Y-%m-%d %H:%M:%S', JSON_EXTRACT_SCALAR(json_payload, '$.data.payload.{{ field }}')) AS {{ field }} {% elif field == 'product_value' %} SAFE_CAST(JSON_EXTRACT_SCALAR(json_payload, '$.data.payload.{{ field }}') AS INT64) AS {{ field }} {% else %} JSON_EXTRACT_SCALAR(json_payload, '$.data.payload.{{ field }}') AS {{ field }} {% endif %}{% if not loop.last %},{% endif %} {% endfor %} FROM {{ source('raw', 'change_events') }} ) SELECT * FROM parsed_data
3. 动态表结构配置
在dbt_project.yml中为该模型添加配置,允许表结构随字段变化动态更新:
models: your_project_name: dynamic_change_table: +materialized: table +allow_refresh: true +schema: target_schema
4. 关键注意事项
- 数据类型适配:根据字段实际业务类型添加转换逻辑,避免存储为字符串类型影响后续分析。
- 容错处理:使用
SAFE_前缀函数(如SAFE_CAST、SAFE_JSON_EXTRACT_SCALAR)避免解析失败导致任务中断。 - 增量优化:如果是增量同步场景,可利用
committimestamp作为增量过滤键,减少全表扫描:WHERE JSON_EXTRACT_SCALAR(json_payload, '$.data.payload.changeeventheader.committimestamp') > (SELECT COALESCE(MAX(committimestamp), 0) FROM {{ this }}) - 字段去重:
changefields_query中使用DISTINCT确保字段列表无重复,避免生成重复的SELECT语句。
内容的提问来源于stack exchange,提问作者Chandra
相关产品推荐
相关产品推荐

