Skip to content

Commit 3e10801

Browse files
authored
fix(clickhouse): strip virtual catalog from view sources (#5940)
Signed-off-by: mday-io <mdaytn@gmail.com>
1 parent 6a563f5 commit 3e10801

2 files changed

Lines changed: 85 additions & 0 deletions

File tree

sqlmesh/core/engine_adapter/clickhouse.py

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -595,6 +595,48 @@ def _create_table(
595595
target_columns_to_types or self.columns(table_name),
596596
)
597597

598+
def create_view(
599+
self,
600+
view_name: TableName,
601+
query_or_df: QueryOrDF,
602+
target_columns_to_types: t.Optional[t.Dict[str, exp.DataType]] = None,
603+
replace: bool = True,
604+
materialized: bool = False,
605+
materialized_properties: t.Optional[t.Dict[str, t.Any]] = None,
606+
table_description: t.Optional[str] = None,
607+
column_descriptions: t.Optional[t.Dict[str, str]] = None,
608+
view_properties: t.Optional[t.Dict[str, exp.Expr]] = None,
609+
source_columns: t.Optional[t.List[str]] = None,
610+
**create_kwargs: t.Any,
611+
) -> None:
612+
if self._default_catalog and isinstance(query_or_df, exp.Query):
613+
from sqlmesh.utils.errors import SQLMeshError
614+
615+
query_or_df = query_or_df.copy()
616+
for table in query_or_df.find_all(exp.Table):
617+
if not table.catalog:
618+
continue
619+
if table.catalog != self._default_catalog:
620+
raise SQLMeshError(
621+
f"{self.dialect} requires that all catalog operations be against a single "
622+
f"catalog: {self._default_catalog}. Provided catalog: {table.catalog}"
623+
)
624+
table.set("catalog", None)
625+
626+
super().create_view(
627+
view_name,
628+
query_or_df,
629+
target_columns_to_types=target_columns_to_types,
630+
replace=replace,
631+
materialized=materialized,
632+
materialized_properties=materialized_properties,
633+
table_description=table_description,
634+
column_descriptions=column_descriptions,
635+
view_properties=view_properties,
636+
source_columns=source_columns,
637+
**create_kwargs,
638+
)
639+
598640
def _strip_virtual_catalog(self, name: "TableName") -> exp.Table:
599641
"""Strip the virtual catalog prefix from a table name if present.
600642

tests/core/engine_adapter/test_clickhouse.py

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1594,3 +1594,46 @@ def test_virtual_catalog_stripped_in_alter_table(make_mocked_engine_adapter: t.C
15941594
assert "mydb" in sql_calls[0]
15951595
assert "my_table" in sql_calls[0]
15961596
assert "ALTER TABLE" in sql_calls[0]
1597+
1598+
1599+
def test_virtual_catalog_stripped_from_create_view_source(
1600+
make_mocked_engine_adapter: t.Callable,
1601+
):
1602+
adapter = make_mocked_engine_adapter(
1603+
ClickhouseEngineAdapter,
1604+
cluster="my_cluster",
1605+
)
1606+
adapter.inject_virtual_catalog("clickhouse_gw")
1607+
query = parse_one("SELECT * FROM __clickhouse_gw__.my_db.my_db__connection_test__1234567890")
1608+
1609+
adapter.create_view(
1610+
"__clickhouse_gw__.my_db.connection_test__dev",
1611+
query,
1612+
)
1613+
1614+
assert to_sql_calls(adapter) == [
1615+
'CREATE OR REPLACE VIEW "my_db"."connection_test__dev" '
1616+
'ON CLUSTER "my_cluster" AS SELECT * FROM '
1617+
'"my_db"."my_db__connection_test__1234567890"'
1618+
]
1619+
assert query.sql() == (
1620+
"SELECT * FROM __clickhouse_gw__.my_db.my_db__connection_test__1234567890"
1621+
)
1622+
1623+
1624+
def test_create_view_source_rejects_unexpected_virtual_catalog(
1625+
make_mocked_engine_adapter: t.Callable,
1626+
):
1627+
from sqlmesh.utils.errors import SQLMeshError
1628+
1629+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1630+
adapter.inject_virtual_catalog("clickhouse_gw")
1631+
1632+
with pytest.raises(
1633+
SQLMeshError,
1634+
match="Provided catalog: unexpected_catalog",
1635+
):
1636+
adapter.create_view(
1637+
"__clickhouse_gw__.my_db.connection_test__dev",
1638+
parse_one("SELECT * FROM unexpected_catalog.my_db.physical_view"),
1639+
)

0 commit comments

Comments
 (0)