Skip to content

In SCD_TYPE_2_BY_TIME models, (null -> non-value) values changes are not tracked properly. #5332

Description

@fbrescia

Issue: SCD_TYPE_2_BY_TIME on BigQuery incorrectly tracks changes of NULL values

SQLMesh version: 0.216.0
Gateway: BigQuery

Description:
When using an SCD_TYPE_2_BY_TIME model kind on a BigQuery backend, NULL values in a source record are being incorrectly filled with non-NULL values from a subsequent update for the same unique_key. This behavior appears to be caused by the use of COALESCE in the generated merge statement, which does not preserve the intended NULL values from the source data.


Steps to reproduce:

  1. Define the following SCD_TYPE_2_BY_TIME model:
MODEL (
  name project.target_model,
  start '2024-01-01',
  columns (
    identifier STRING,
    a_nullable_value STRING,
    another_nullable_value STRING,
    updated_at TIMESTAMP
  ),
  kind SCD_TYPE_2_BY_TIME (
    unique_key identifier,
    updated_at_name updated_at,
    time_data_type TIMESTAMP,
    batch_size 1,
    forward_only false
  ),
  partitioned_by TIMESTAMP_TRUNC(valid_from, DAY),
  cron '@daily'
);

SELECT
  identifier,
  a_nullable_value,
  another_nullable_value,
  updated_at
FROM source-data.source_dataset.source_table
WHERE
  _PARTITIONTIME BETWEEN @start_ds AND @end_ds
  1. Provide the following source data in source-data.source_dataset.source_table:
identifier  a_nullable_value  another_nullable_value  updated_at                      _PARTITIONTIME                
'aaa'       (null)            'C'                     2024-02-03 17:25:04.000000 UTC  2024-02-03 00:00:00.000000 UTC
'aaa'       'A'               'C'                     2024-08-07 05:45:58.000000 UTC  2024-08-07 00:00:00.000000 UTC
'aaa'       'B'               'D'                     2025-02-06 14:41:54.000000 UTC  2025-02-06 00:00:00.000000 UTC
  1. Run a sqlmesh plan to create and populate the target model.

Expected behavior:

The initial record from 2024-02-03 should retain its NULL value for the a_nullable_value column.
The resulting project.target_model table should contain the following entries:

identifier  a_nullable_value  another_nullable_value  updated_at                      valid_from                      valid_to                      
'aaa'       (null)            'C'                     2024-02-03 17:25:04.000000 UTC  2024-02-03 17:25:04.000000 UTC  2024-08-07 05:45:58.000000 UTC
'aaa'       'A'               'C'                     2024-08-07 05:45:58.000000 UTC  2024-08-07 05:45:58.000000 UTC  2025-02-06 14:41:54.000000 UTC
'aaa'       'B'               'D'                     2025-02-06 14:41:54.000000 UTC  2025-02-06 14:41:54.000000 UTC  (null)                        

Actual Behavior:

The NULL value in the a_nullable_value column for the first record is replaced by the value 'A' from the subsequent record.
The table is populated with the following incorrect data:

identifier  a_nullable_value  another_nullable_value  updated_at                      valid_from                      valid_to                      
'aaa'       'A'               'C'                     2024-02-03 17:25:04.000000 UTC  2024-02-03 17:25:04.000000 UTC  2024-08-07 05:45:58.000000 UTC
'aaa'       'A'               'C'                     2024-08-07 05:45:58.000000 UTC  2024-08-07 05:45:58.000000 UTC  2025-02-06 14:41:54.000000 UTC
'aaa'       'B'               'D'                     2025-02-06 14:41:54.000000 UTC  2025-02-06 14:41:54.000000 UTC  (null)                        

Possible Cause:

The issue likely stems from the generated CREATE OR REPLACE TABLE statement, which uses COALESCE on all columns.
This logic incorrectly backfills NULLs with values from later records during the join operation.

Relevant Query Snippet:

  SELECT
    COALESCE(`joined`.`t_identifier`, `joined`.`identifier`) AS `identifier`,
    COALESCE(`joined`.`t_a_nullable_value`, `joined`.`a_nullable_value`) AS `a_nullable_value`,
    COALESCE(`joined`.`t_another_nullable_value`, `joined`.`another_nullable_value`) AS `another_nullable_value`,
    COALESCE(`joined`.`t_updated_at`, `joined`.`updated_at`) AS `updated_at`,
    CASE
      WHEN `t_valid_from` IS NULL AND NOT `latest_deleted`.`_exists` IS NULL
        THEN
          CASE
            WHEN `latest_deleted`.`valid_to` > `updated_at`
              THEN `latest_deleted`.`valid_to`
            ELSE `updated_at`
          END
      WHEN `t_valid_from` IS NULL
        THEN `updated_at`
      ELSE `t_valid_from`
    END AS `valid_from`,
    CASE
      WHEN `joined`.`updated_at` > `joined`.`t_updated_at`
        THEN `joined`.`updated_at`
      ELSE `t_valid_to`
    END AS `valid_to`
  FROM `joined`
  LEFT JOIN `latest_deleted` ON `joined`.`identifier` = `latest_deleted`.`_key0`

Activity

  1. esatiukov commented on Mar 12, 2026

    @esatiukov

    +1 on Postgres engine

  2. StuffbyYuki commented on Mar 24, 2026

    @StuffbyYuki
    Collaborator

    Do we have any updates on this bug?

  3. ionnich commented on Oct 2, 2026

    @ionnich

    Still reproduces on SQLMesh 0.236.2, and not only on BigQuery or BY_TIME: SCD_TYPE_2_BY_COLUMN (with check_columns) and SCD_TYPE_2_BY_TIME share EngineAdapter._scd_type_2, and both lose the NULL. We reproduced it on Spark (Sail 0.7.2, Iceberg).

    Root cause. In sqlmesh/core/engine_adapter/base.py, _scd_type_2 builds the updated_rows CTE. That CTE re-emits every current target row (the version being closed, or the one kept), and it takes each unmanaged column as

    exp.func(
        "COALESCE",
        exp.column(prefixed_unmanaged_columns[i].this, table="joined"),  # joined.t_<col>, the target row
        exp.column(col, table="joined"),                                 # joined.<col>, the new source row
    ).as_(col)

    Suppose the target row exists and its value is NULL. COALESCE then falls through to the source row's value, so the version being closed shows the value that only the next version had. inserted_rows still opens the new version correctly, so version counts and current rows look right. Only the history is wrong, and SCD2 kinds are forward-only and cannot be restated, so the wrong history stays.

    Minimal repro (BY_COLUMN, check_columns [v], execution_time_as_valid_from true):

    load source (k, v)
    2026-09-29 (n, NULL)
    2026-09-30 (n, 'A')
    2026-10-01 (n, 'B')

    Expected history:

    n  NULL  2026-09-29  2026-09-30
    n  'A'   2026-09-30  2026-10-01
    n  'B'   2026-10-01  NULL
    

    Actual: the first row reads 'A' (n, 'A', 2026-09-29, 2026-09-30). The value→NULL direction (x then NULL) and unchanged keys come out correct. BY_TIME with updated_at shows the same NULL→'A' overwrite.

    On a real dimension (93k symbols, 4 daily loads), 283 of the 574 versions closed in one run were identical to their successor: FMP had filled in a NULL cik. That is history which never existed.

    Fix. Prefer the target value whenever a target row exists. _scd_type_2 already uses t_<valid_from> IS NULL as its "no target row" test (valid_from_case_stmt in both branches), so use that:

    --- a/sqlmesh/core/engine_adapter/base.py
    +++ b/sqlmesh/core/engine_adapter/base.py
    @@ def _scd_type_2(
                     .with_(
                         "updated_rows",
                         exp.select(
                             *(
    -                            exp.func(
    -                                "COALESCE",
    -                                exp.column(prefixed_unmanaged_columns[i].this, table="joined"),
    -                                exp.column(col, table="joined"),
    -                            ).as_(col)
    +                            exp.Case()
    +                            .when(
    +                                exp.column(prefixed_valid_from_col.this, table="joined")
    +                                .is_(exp.Null())
    +                                .not_(),
    +                                exp.column(prefixed_unmanaged_columns[i].this, table="joined"),
    +                            )
    +                            .else_(exp.column(col, table="joined"))
    +                            .as_(col)
                                 for i, col in enumerate(unmanaged_columns_to_types)
                             ),

    Patch description.

    fix(scd2): keep NULLs in closed versions

    updated_rows coalesced each target column with the source column, so a NULL in the version being closed (or kept) was replaced by the new source value. Now the target value is taken whenever a target row exists (t_<valid_from> is not NULL, the same test valid_from already uses), and the source value only for new keys. This applies to BY_TIME and BY_COLUMN. Rows for new keys, deleted keys and unchanged keys are unaffected. Behaviour change for existing tables: from the next run, closed versions keep their NULLs. History already written is not rewritten.

    Tests: add a NULL→value→value key to the SCD2 adapter tests and assert that the closed version keeps NULL. Any expected-SQL assertion containing COALESCE("joined"."t_<col>", "joined"."<col>") needs the CASE form (not checked against the upstream test suite).

    The existence check uses t_<valid_from>, not t_<unique_key>: a target row whose key part is NULL never matches the join, and a key-based check would then emit that row with all-NULL source values.

    We run the same rewrite as an engine-adapter override (ensure_nulls_for_unmatched_after_join, applied to the finished SCD2 query) until this is released.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    BugSomething isn't working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions