diff --git a/sqlmesh/core/engine_adapter/base.py b/sqlmesh/core/engine_adapter/base.py index 930cdf7cd4..48aac3b56e 100644 --- a/sqlmesh/core/engine_adapter/base.py +++ b/sqlmesh/core/engine_adapter/base.py @@ -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 diff --git a/sqlmesh/core/engine_adapter/redshift.py b/sqlmesh/core/engine_adapter/redshift.py index 8f2f098d08..1039acc4ea 100644 --- a/sqlmesh/core/engine_adapter/redshift.py +++ b/sqlmesh/core/engine_adapter/redshift.py @@ -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 diff --git a/sqlmesh/core/snapshot/evaluator.py b/sqlmesh/core/snapshot/evaluator.py index 9d9ecf7672..976988ba93 100644 --- a/sqlmesh/core/snapshot/evaluator.py +++ b/sqlmesh/core/snapshot/evaluator.py @@ -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 diff --git a/tests/core/test_snapshot_evaluator.py b/tests/core/test_snapshot_evaluator.py index a9ca86496a..c2702f4a1a 100644 --- a/tests/core/test_snapshot_evaluator.py +++ b/tests/core/test_snapshot_evaluator.py @@ -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, @@ -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: @@ -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 ):