Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions sqlmesh/core/engine_adapter/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ class EngineAdapter:
SCHEMA_DIFFER_KWARGS: t.Dict[str, t.Any] = {}
SUPPORTS_TUPLE_IN = True
HAS_VIEW_BINDING = False
RECREATE_VIEW_ON_EVALUATION = True
RECREATE_MATERIALIZED_VIEW_ON_EVALUATION = True
SUPPORTS_REPLACE_TABLE = True
SUPPORTS_GRANTS = False
Expand Down
4 changes: 4 additions & 0 deletions sqlmesh/core/engine_adapter/redshift.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,10 @@ class RedshiftEngineAdapter(
# Redshift doesn't support comments for VIEWs WITH NO SCHEMA BINDING (which we always use)
COMMENT_CREATION_VIEW = CommentCreationView.UNSUPPORTED
SUPPORTS_REPLACE_TABLE = False
# Replacing a view (DROP + CREATE) gives it a new OID, which breaks concurrent queries that read through it.
# Views WITH NO SCHEMA BINDING resolve their dependencies at query time, so an existing view only needs to be
# recreated on the first insert of a snapshot version.
RECREATE_VIEW_ON_EVALUATION = False
SUPPORTS_GRANTS = True
SUPPORTS_MULTIPLE_GRANT_PRINCIPALS = True

Expand Down
10 changes: 10 additions & 0 deletions sqlmesh/core/snapshot/evaluator.py
Original file line number Diff line number Diff line change
Expand Up @@ -2750,6 +2750,16 @@ def insert(
):
must_recreate_view = False

# Some engines (e.g. Redshift) replace views via DROP + CREATE, which breaks concurrent queries that read
# through the view. For those, an existing view must not be recreated on routine evaluation; only on the
# first insert, which also covers rebuilds forced by `should_force_rebuild`.
if (
not is_materialized_view
and not is_first_insert
and not self.adapter.RECREATE_VIEW_ON_EVALUATION
):
must_recreate_view = False

if self.adapter.table_exists(table_name) and not must_recreate_view:
logger.info("Skipping creation of the view '%s'", table_name)
return
Expand Down
62 changes: 61 additions & 1 deletion tests/core/test_snapshot_evaluator.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,13 @@
from sqlmesh.core.audit import ModelAudit, StandaloneAudit
from sqlmesh.core import dialect as d
from sqlmesh.core.dialect import schema_, to_schema
from sqlmesh.core.engine_adapter import EngineAdapter, create_engine_adapter, BigQueryEngineAdapter
from sqlmesh.core.engine_adapter import (
EngineAdapter,
create_engine_adapter,
BigQueryEngineAdapter,
RedshiftEngineAdapter,
SnowflakeEngineAdapter,
)
from sqlmesh.core.engine_adapter.base import MERGE_SOURCE_ALIAS, MERGE_TARGET_ALIAS
from sqlmesh.core.engine_adapter.shared import (
DataObject,
Expand Down Expand Up @@ -81,6 +87,7 @@
)
from sqlmesh.utils.metaprogramming import Executable
from sqlmesh.utils.pydantic import list_of_fields_validator
from tests.core.engine_adapter import to_sql_calls


if t.TYPE_CHECKING:
Expand Down Expand Up @@ -764,6 +771,59 @@ def test_evaluate_materialized_view_not_recreated_on_evaluation(
assert adapter_mock.create_view.call_count == expected_create_view_calls


@pytest.mark.parametrize(
"adapter_cls, materialized, has_intervals, expect_recreate",
[
# Existing view with intervals -> routine evaluation: do NOT recreate on Redshift
(RedshiftEngineAdapter, False, True, False),
# Existing view without intervals -> first insert (e.g. `should_force_rebuild`): recreate it
(RedshiftEngineAdapter, False, False, True),
# Materialized views are still recreated on every evaluation
(RedshiftEngineAdapter, True, True, True),
# Other engines without view binding still recreate existing views on every evaluation
(SnowflakeEngineAdapter, False, True, True),
(SnowflakeEngineAdapter, False, False, True),
(SnowflakeEngineAdapter, True, True, True),
],
)
def test_evaluate_existing_view_recreation(
mocker: MockerFixture,
make_mocked_engine_adapter,
make_snapshot,
adapter_cls: t.Type[EngineAdapter],
materialized: bool,
has_intervals: bool,
expect_recreate: bool,
):
adapter = make_mocked_engine_adapter(adapter_cls)
adapter.with_settings = lambda **kwargs: adapter # type: ignore
mocker.patch.object(adapter, "table_exists", return_value=True)
evaluator = SnapshotEvaluator(adapter)

model = SqlModel(
name="test_schema.test_model",
kind=ViewKind(materialized=materialized),
query=parse_one("SELECT a FROM tbl"),
)
snapshot = make_snapshot(model)
snapshot.categorize_as(SnapshotChangeCategory.BREAKING)
if has_intervals:
snapshot.add_interval("2023-01-01", "2023-01-01")

evaluator.evaluate(
snapshot,
start="2020-01-01",
end="2020-01-02",
execution_time="2020-01-02",
snapshots={},
)

view_ddl = [
sql for sql in to_sql_calls(adapter) if sql.startswith(("CREATE", "DROP")) and "VIEW" in sql
]
assert bool(view_ddl) == expect_recreate


def test_evaluate_materialized_view_with_partitioned_by_cluster_by(
mocker: MockerFixture, adapter_mock, make_snapshot
):
Expand Down