Skip to content

Commit bf1d818

Browse files
Arthur031221ionnich
andcommitted
fix(scd2): preserve NULL values in historical rows
Select target values using an explicit row-existence flag carried into the SCD2 join. Choose a flag alias that does not overlap source or prefixed target columns, including case-insensitive names. Co-authored-by: Nicholas <71755971+ionnich@users.noreply.github.com> Signed-off-by: Arthur031221 <levi74108520963@gmail.com>
1 parent 271b855 commit bf1d818

5 files changed

Lines changed: 504 additions & 49 deletions

File tree

‎sqlmesh/core/engine_adapter/base.py‎

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
from sqlglot import Dialect, exp
2020
from sqlglot.errors import ErrorLevel
21-
from sqlglot.helper import ensure_list, seq_get
21+
from sqlglot.helper import ensure_list, find_new_name, seq_get
2222
from sqlglot.optimizer.qualify_columns import quote_identifiers
2323

2424
from sqlmesh.core.dialect import (
@@ -2073,6 +2073,11 @@ def remove_managed_columns(
20732073
prefixed_col = exp.column(column).copy()
20742074
prefixed_col.this.set("this", f"t_{prefixed_col.name}")
20752075
prefixed_unmanaged_columns.append(prefixed_col)
2076+
target_exists_column = find_new_name(
2077+
{col.name.lower() for col in prefixed_columns_to_types}
2078+
| {col.lower() for col in unmanaged_columns_to_types},
2079+
"t__exists",
2080+
)
20762081
query = (
20772082
exp.Select() # type: ignore
20782083
.select(*table_columns)
@@ -2141,6 +2146,7 @@ def remove_managed_columns(
21412146
"joined",
21422147
exp.select(
21432148
exp.column("_exists", table="source").as_("_exists"),
2149+
exp.column("_exists", table="latest").as_(target_exists_column),
21442150
*(
21452151
exp.column(col, table="latest").as_(prefixed_columns_to_types[i].this)
21462152
for i, col in enumerate(target_columns_to_types)
@@ -2164,6 +2170,7 @@ def remove_managed_columns(
21642170
.union(
21652171
exp.select(
21662172
exp.column("_exists", table="source").as_("_exists"),
2173+
exp.column("_exists", table="latest").as_(target_exists_column),
21672174
*(
21682175
exp.column(col, table="latest").as_(
21692176
prefixed_columns_to_types[i].this
@@ -2195,11 +2202,15 @@ def remove_managed_columns(
21952202
"updated_rows",
21962203
exp.select(
21972204
*(
2198-
exp.func(
2199-
"COALESCE",
2205+
exp.Case()
2206+
.when(
2207+
exp.column(target_exists_column, table="joined")
2208+
.is_(exp.Null())
2209+
.not_(),
22002210
exp.column(prefixed_unmanaged_columns[i].this, table="joined"),
2201-
exp.column(col, table="joined"),
2202-
).as_(col)
2211+
)
2212+
.else_(exp.column(col, table="joined"))
2213+
.as_(col)
22032214
for i, col in enumerate(unmanaged_columns_to_types)
22042215
),
22052216
valid_from_case_stmt,

0 commit comments

Comments
 (0)