通过dbt宏遍历BigQuery结构体字段生成新列的问题
在dbt中处理BigQuery STRUCT列,生成动态原因数组
问题诊断
你原代码的核心问题是混淆了dbt的编译时逻辑和运行时SQL执行逻辑:
dbt_utils.get_query_results_as_dict是在dbt编译SQL阶段拉取全表数据,遍历的是整个结果集的id和rating_record集合,而非每行单独处理,导致生成的drop_reasons和bump_reasons是所有行的汇总值,无法对应到单个id。- Jinja循环变量是编译时生成的,无法关联到SQL运行时的行级数据,最终生成的SQL逻辑完全不符合预期。
解决方案
针对BigQuery的STRUCT类型,推荐两种更高效的实现方式:
方法1:用BigQuery原生函数动态处理STRUCT(推荐)
利用BigQuery的STRUCT_TO_ARRAY和UNNEST函数,在SQL层面直接处理行级STRUCT数据,无需依赖dbt宏遍历数据:
WITH struct_key_values AS ( SELECT id, -- 将STRUCT转换为键值对数组,每个元素包含字段名和对应值 UNNEST(STRUCT_TO_ARRAY(rating_record)) AS kv_pair FROM {{ ref('rating') }} ) SELECT id, -- 聚合所有值为true的_drop后缀字段名 ARRAY_AGG(DISTINCT kv_pair.key) FILTER (WHERE kv_pair.key LIKE '%_drop' AND kv_pair.value) AS drop_reasons, -- 聚合所有值为true的_bump后缀字段名 ARRAY_AGG(DISTINCT kv_pair.key) FILTER (WHERE kv_pair.key LIKE '%_bump' AND kv_pair.value) AS bump_reasons FROM struct_key_values GROUP BY id
这种方法的优势是无需提前知道STRUCT的具体字段,即使后续STRUCT新增字段,代码也能自动适配。
方法2:用dbt宏获取STRUCT字段生成静态逻辑
如果需要明确基于表结构生成SQL(比如字段固定),可以用dbt_utils.get_columns_by_table获取STRUCT的子字段,动态生成判断逻辑:
{% set struct_subfields = dbt_utils.get_columns_by_table(ref('rating'), 'rating_record') %} SELECT id, -- 生成drop原因数组:过滤出值为true的_drop字段 [ {% for field in struct_subfields %} {% if field.name.endswith('_drop') %} CASE WHEN rating_record.{{ field.name }} THEN '{{ field.name }}' END{% if not loop.last %},{% endif %} {% endif %} {% endfor %} ] FILTER (WHERE VALUE IS NOT NULL) AS drop_reasons, -- 生成bump原因数组:过滤出值为true的_bump字段 [ {% for field in struct_subfields %} {% if field.name.endswith('_bump') %} CASE WHEN rating_record.{{ field.name }} THEN '{{ field.name }}' END{% if not loop.last %},{% endif %} {% endif %} {% endfor %} ] FILTER (WHERE VALUE IS NOT NULL) AS bump_reasons FROM {{ ref('rating') }}
这种方法会在dbt编译阶段获取STRUCT的字段列表,生成对应的CASE WHEN逻辑,最终SQL会更简洁。
原代码问题总结
你之前的思路错误地试图在dbt编译阶段处理运行时的行级数据,这违背了dbt的工作机制:Jinja模板负责生成SQL,而SQL的行级计算是在BigQuery执行阶段完成的,两者不能混淆使用。
内容的提问来源于stack exchange,提问作者Faisal Maqbool
相关产品推荐
相关产品推荐

