增量Merge仅同步公共列:source_columns配置失效问题排查
增量Merge场景下源表与目标表Schema不一致的处理问题
问题定位
核心问题是dbt的config参数无法传递复杂对象:
adapter.get_columns_in_relation()返回的是dbt内部的Column对象列表,属于带属性的复杂结构,而dbt的config系统仅支持可JSON序列化的简单类型(字符串、数字、普通列表/字典)。- 当你把
src_cols(Column对象列表)传入source_columns配置时,这些对象会被强制序列化为字符串(比如<Column(Entity_ID)>),导致宏中config.get('source_columns')拿到的不是预期的列名结构。
临时修复方案
先将源表列提取为字符串列表,再传入config:
Model代码调整
{%- set src_cols = adapter.get_columns_in_relation(ref('pre_Dim_Entities_Client')) | map(attribute='name') | list -%} {{ config( materialized='incremental', unique_key='Entity_ID', source_columns = src_cols ) }} SELECT * FROM {{ ref ('pre_Dim_Entities_Client')}}
宏代码调整
{% macro default__get_merge_sql(target, source, unique_key, dest_columns, predicates) -%} {%- set predicates = [] if predicates is none else [] + predicates -%} {%- set src_col_names = config.get('source_columns') -%} {{- print(src_col_names) -}} {%- set dest_cols_csv = get_quoted_csv(src_col_names if src_col_names else dest_columns | map(attribute="name")) -%} {%- set update_columns = config.get('merge_update_columns', default = src_col_names | map('quote') | list if src_col_names else dest_columns | map(attribute="quoted") | list) -%} {%- set sql_header = config.get('sql_header', none) -%} -- 此处补充完整的Merge SQL逻辑(按原有逻辑实现剩余部分) {% endmacro %}
更优实现方案
方案1:宏内自动计算公共列(推荐)
无需在model中传递列信息,宏内自动获取源表与目标表的公共列,减少冗余配置:
宏代码
{% macro default__get_merge_sql(target, source, unique_key, dest_columns, predicates) -%} {%- set predicates = [] if predicates is none else [] + predicates -%} -- 从model的config中读取源表关系,获取源表列名列表 {%- set src_relation = ref(config.get('source_relation')) -%} {%- set src_col_names = adapter.get_columns_in_relation(src_relation) | map(attribute='name') | list -%} -- 获取目标表列名列表 {%- set dest_col_names = dest_columns | map(attribute='name') | list -%} -- 计算两者的公共列 {%- set common_cols = src_col_names | intersect(dest_col_names) -%} {%- set dest_cols_csv = get_quoted_csv(common_cols) -%} {%- set update_columns = common_cols | map('quote') | list -%} {%- set sql_header = config.get('sql_header', none) -%} -- 生成完整Merge SQL {{ sql_header if sql_header is not none }} merge into {{ target }} as DBT_INTERNAL_DEST using {{ source }} as DBT_INTERNAL_SOURCE on DBT_INTERNAL_DEST.{{ unique_key }} = DBT_INTERNAL_SOURCE.{{ unique_key }} {% if predicates %} {% for predicate in predicates %} and {{ predicate }} {% endfor %} {% endif %} when matched then update set {% for col in update_columns %} {{ col }} = DBT_INTERNAL_SOURCE.{{ col }} {% if not loop.last %},{% endif %} {% endfor %} when not matched then insert ({{ dest_cols_csv }}) values ({{ get_quoted_csv(common_cols, prefix='DBT_INTERNAL_SOURCE.') }}) {% endmacro %}
Model代码
{{ config( materialized='incremental', unique_key='Entity_ID', source_relation='pre_Dim_Entities_Client' -- 仅需指定源表名 ) }} SELECT * FROM {{ ref ('pre_Dim_Entities_Client')}}
方案2:Model内提前计算公共列
在model中预先计算源表与目标表的公共列,仅选择需要的列进行处理,避免冗余数据:
Model代码
{%- set src_relation = ref('pre_Dim_Entities_Client') -%} {%- set src_cols = adapter.get_columns_in_relation(src_relation) | map(attribute='name') | list -%} {%- set dest_cols = adapter.get_columns_in_relation(this) | map(attribute='name') | list -%} {%- set common_cols = src_cols | intersect(dest_cols) -%} {{ config( materialized='incremental', unique_key='Entity_ID', common_columns = common_cols ) }} -- 仅选择公共列,避免不必要的列加载 SELECT {{ get_quoted_csv(common_cols) }} FROM {{ src_relation }}
宏代码
{% macro default__get_merge_sql(target, source, unique_key, dest_columns, predicates) -%} {%- set predicates = [] if predicates is none else [] + predicates -%} {%- set common_cols = config.get('common_columns') -%} {%- set dest_cols_csv = get_quoted_csv(common_cols) -%} {%- set update_columns = common_cols | map('quote') | list -%} {%- set sql_header = config.get('sql_header', none) -%} -- 此处补充完整Merge SQL逻辑 {% endmacro %}
内容的提问来源于stack exchange,提问作者Dina Kleper
相关产品推荐
相关产品推荐

