Skip to content

Commit 00d46ed

Browse files
authored
fix(clickhouse): strip virtual catalogs from executed queries and inserts (#6036)
Signed-off-by: mday-io <mdaytn@gmail.com>
1 parent 3e0a5ca commit 00d46ed

2 files changed

Lines changed: 163 additions & 16 deletions

File tree

‎sqlmesh/core/engine_adapter/clickhouse.py‎

Lines changed: 37 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,21 @@ def inject_virtual_catalog(self, gateway: str) -> None:
6161
configured = self._extra_config.get("virtual_catalog")
6262
self._default_catalog = f"__{gateway}__" if configured is None else configured
6363

64+
def _to_sql(self, expression: exp.Expr, quote: bool = True, **kwargs: t.Any) -> str:
65+
"""Render SQL without the virtual catalog, which ClickHouse does not support."""
66+
if self._default_catalog and any(self._virtual_catalog_references(expression)):
67+
expression = expression.copy()
68+
for reference in list(self._virtual_catalog_references(expression)):
69+
reference.set("catalog", None)
70+
return super()._to_sql(expression, quote=quote, **kwargs)
71+
72+
def _virtual_catalog_references(self, expression: exp.Expr) -> t.Iterator[exp.Expr]:
73+
return (
74+
reference
75+
for reference in expression.find_all(exp.Table, exp.Column)
76+
if reference.text("catalog") == self._default_catalog
77+
)
78+
6479
@property
6580
def engine_run_mode(self) -> EngineRunMode:
6681
if self._extra_config.get("cloud_mode"):
@@ -502,8 +517,14 @@ def _create_table_like(
502517
**kwargs: t.Any,
503518
) -> None:
504519
"""Create table with identical structure as source table"""
520+
target_table_sql = self._strip_virtual_catalog(target_table_name).sql(
521+
dialect=self.dialect, identify=True
522+
)
523+
source_table_sql = self._strip_virtual_catalog(source_table_name).sql(
524+
dialect=self.dialect, identify=True
525+
)
505526
self.execute(
506-
f"CREATE TABLE {target_table_name}{self._on_cluster_sql()} AS {source_table_name}"
527+
f"CREATE TABLE {target_table_sql}{self._on_cluster_sql()} AS {source_table_sql}"
507528
)
508529

509530
def _get_partition_ids(
@@ -648,7 +669,7 @@ def _strip_virtual_catalog(self, name: "TableName") -> exp.Table:
648669
SQL is sent to the wire, since ClickHouse only supports a two-level
649670
``[database].[table]`` naming scheme.
650671
"""
651-
table = exp.to_table(name)
672+
table = exp.to_table(name, dialect=self.dialect)
652673
if self._default_catalog and table.catalog == self._default_catalog:
653674
table.set("catalog", None)
654675
return table
@@ -660,8 +681,12 @@ def _exchange_tables(
660681
) -> None:
661682
from clickhouse_connect.driver.exceptions import DatabaseError # type: ignore
662683

663-
old_table_sql = exp.to_table(old_table_name).sql(dialect=self.dialect, identify=True)
664-
new_table_sql = exp.to_table(new_table_name).sql(dialect=self.dialect, identify=True)
684+
old_table_sql = self._strip_virtual_catalog(old_table_name).sql(
685+
dialect=self.dialect, identify=True
686+
)
687+
new_table_sql = self._strip_virtual_catalog(new_table_name).sql(
688+
dialect=self.dialect, identify=True
689+
)
665690

666691
try:
667692
self.execute(
@@ -685,8 +710,12 @@ def _rename_table(
685710
old_table_name: TableName,
686711
new_table_name: TableName,
687712
) -> None:
688-
old_table_sql = exp.to_table(old_table_name).sql(dialect=self.dialect, identify=True)
689-
new_table_sql = exp.to_table(new_table_name).sql(dialect=self.dialect, identify=True)
713+
old_table_sql = self._strip_virtual_catalog(old_table_name).sql(
714+
dialect=self.dialect, identify=True
715+
)
716+
new_table_sql = self._strip_virtual_catalog(new_table_name).sql(
717+
dialect=self.dialect, identify=True
718+
)
690719

691720
self.execute(f"RENAME TABLE {old_table_sql} TO {new_table_sql}{self._on_cluster_sql()}")
692721

@@ -974,7 +1003,7 @@ def _build_view_properties_exp(
9741003
def _build_create_comment_table_exp(
9751004
self, table: exp.Table, table_comment: str, table_kind: str, **kwargs: t.Any
9761005
) -> exp.Comment | str:
977-
table_sql = table.sql(dialect=self.dialect, identify=True)
1006+
table_sql = self._strip_virtual_catalog(table).sql(dialect=self.dialect, identify=True)
9781007

9791008
truncated_comment = self._truncate_table_comment(table_comment)
9801009
comment_sql = exp.Literal.string(truncated_comment).sql(dialect=self.dialect)
@@ -989,7 +1018,7 @@ def _build_create_comment_column_exp(
9891018
table_kind: str = "TABLE",
9901019
**kwargs: t.Any,
9911020
) -> exp.Comment | str:
992-
table_sql = table.sql(dialect=self.dialect, identify=True)
1021+
table_sql = self._strip_virtual_catalog(table).sql(dialect=self.dialect, identify=True)
9931022
column_sql = exp.to_column(column_name).sql(dialect=self.dialect, identify=True)
9941023

9951024
truncated_comment = self._truncate_table_comment(column_comment)

‎tests/core/engine_adapter/test_clickhouse.py‎

Lines changed: 126 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1066,7 +1066,7 @@ def test_insert_overwrite_by_condition_replace_partitioned(
10661066
)
10671067

10681068
assert to_sql_calls(adapter) == [
1069-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1069+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
10701070
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT * FROM "__temp_new_records_abcd") AS "_subquery"',
10711071
'EXCHANGE TABLES "__temp_existing_records_abcd" AND "__temp_target_abcd"',
10721072
'DROP TABLE IF EXISTS "__temp_target_abcd"',
@@ -1104,7 +1104,7 @@ def test_insert_overwrite_by_condition_replace(
11041104
)
11051105

11061106
to_sql_calls(adapter) == [
1107-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1107+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
11081108
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT * FROM "__temp_new_records_abcd") AS "_subquery"',
11091109
'EXCHANGE TABLES "__temp_existing_records_abcd" AND "__temp_target_abcd"',
11101110
'DROP TABLE IF EXISTS "__temp_target_abcd"',
@@ -1153,7 +1153,7 @@ def test_insert_overwrite_by_condition_where_partitioned(
11531153
)
11541154

11551155
to_sql_calls(adapter) == [
1156-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1156+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
11571157
"""INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT * FROM "__temp_new_records_abcd") AS "_subquery" WHERE "ds" BETWEEN '2024-02-15' AND '2024-04-30'""",
11581158
"""CREATE TABLE IF NOT EXISTS "__temp_target_abcd" ENGINE=MergeTree ORDER BY () AS SELECT DISTINCT "partition_id" FROM (SELECT "_partition_id" AS "partition_id" FROM "__temp_existing_records_abcd" WHERE "ds" BETWEEN '2024-02-15' AND '2024-04-30' UNION DISTINCT SELECT "_partition_id" AS "partition_id" FROM "__temp_target_abcd") AS "_affected_partitions\"""",
11591159
"""INSERT INTO "__temp_target_abcd" SELECT "id", "ds" FROM "__temp_existing_records_abcd" WHERE NOT ("ds" BETWEEN '2024-02-15' AND '2024-04-30') AND "_partition_id" IN (SELECT "partition_id" FROM "__temp_target_abcd")""",
@@ -1204,12 +1204,12 @@ def test_insert_overwrite_by_condition_by_key(
12041204
)
12051205

12061206
to_sql_calls(adapter) == [
1207-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1207+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
12081208
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT DISTINCT ON ("id") * FROM "__temp_new_records_abcd") AS "_subquery"',
12091209
'INSERT INTO "__temp_target_abcd" SELECT "id", "ds" FROM "__temp_existing_records_abcd" WHERE NOT ("id" IN (SELECT "id" FROM "__temp_target_abcd"))',
12101210
'EXCHANGE TABLES "__temp_existing_records_abcd" AND "__temp_target_abcd"',
12111211
'DROP TABLE IF EXISTS "__temp_target_abcd"',
1212-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1212+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
12131213
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT * FROM "__temp_new_records_abcd") AS "_subquery"',
12141214
'INSERT INTO "__temp_target_abcd" SELECT "id", "ds" FROM "__temp_existing_records_abcd" WHERE NOT ("id" IN (SELECT "id" FROM "__temp_target_abcd"))',
12151215
'EXCHANGE TABLES "__temp_existing_records_abcd" AND "__temp_target_abcd"',
@@ -1267,13 +1267,13 @@ def test_insert_overwrite_by_condition_by_key_partitioned(
12671267
)
12681268

12691269
to_sql_calls(adapter) == [
1270-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1270+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
12711271
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT DISTINCT ON ("id") * FROM "__temp_new_records_abcd") AS "_subquery"',
12721272
'CREATE TABLE IF NOT EXISTS "__temp_target_abcd" ENGINE=MergeTree ORDER BY () AS SELECT DISTINCT "partition_id" FROM (SELECT "_partition_id" AS "partition_id" FROM "__temp_existing_records_abcd" WHERE "id" IN (SELECT "id" FROM "__temp_target_abcd") UNION DISTINCT SELECT "_partition_id" AS "partition_id" FROM "__temp_target_abcd") AS "_affected_partitions"',
12731273
'INSERT INTO "__temp_target_abcd" SELECT "id", "ds" FROM "__temp_existing_records_abcd" WHERE NOT ("id" IN (SELECT "id" FROM "__temp_target_abcd")) AND "_partition_id" IN (SELECT "partition_id" FROM "__temp_target_abcd")',
12741274
"""ALTER TABLE "__temp_existing_records_abcd" REPLACE PARTITION ID '2' FROM "__temp_target_abcd", REPLACE PARTITION ID '1' FROM "__temp_target_abcd", REPLACE PARTITION ID '4' FROM "__temp_target_abcd", DROP PARTITION ID '3'""",
12751275
'DROP TABLE IF EXISTS "__temp_target_abcd"',
1276-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1276+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
12771277
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT * FROM "__temp_new_records_abcd") AS "_subquery"',
12781278
'CREATE TABLE IF NOT EXISTS "__temp_target_abcd" ENGINE=MergeTree ORDER BY () AS SELECT DISTINCT "partition_id" FROM (SELECT "_partition_id" AS "partition_id" FROM "__temp_existing_records_abcd" WHERE "id" IN (SELECT "id" FROM "__temp_target_abcd") UNION DISTINCT SELECT "_partition_id" AS "partition_id" FROM "__temp_target_abcd") AS "_affected_partitions"',
12791279
'INSERT INTO "__temp_target_abcd" SELECT "id", "ds" FROM "__temp_existing_records_abcd" WHERE NOT ("id" IN (SELECT "id" FROM "__temp_target_abcd")) AND "_partition_id" IN (SELECT "partition_id" FROM "__temp_target_abcd")',
@@ -1316,7 +1316,7 @@ def test_insert_overwrite_by_condition_inc_by_partition(
13161316
)
13171317

13181318
to_sql_calls(adapter) == [
1319-
"CREATE TABLE __temp_target_abcd AS __temp_existing_records_abcd",
1319+
'CREATE TABLE "__temp_target_abcd" AS "__temp_existing_records_abcd"',
13201320
'INSERT INTO "__temp_target_abcd" ("id", "ds") SELECT "id", "ds" FROM (SELECT * FROM "__temp_new_records_abcd") AS "_subquery"',
13211321
"""ALTER TABLE "__temp_existing_records_abcd" REPLACE PARTITION ID '1' FROM "__temp_target_abcd", REPLACE PARTITION ID '2' FROM "__temp_target_abcd", REPLACE PARTITION ID '4' FROM "__temp_target_abcd\"""",
13221322
'DROP TABLE IF EXISTS "__temp_target_abcd"',
@@ -1596,6 +1596,124 @@ def test_virtual_catalog_stripped_in_alter_table(make_mocked_engine_adapter: t.C
15961596
assert "ALTER TABLE" in sql_calls[0]
15971597

15981598

1599+
@pytest.mark.parametrize(
1600+
"query_sql, expected_sql",
1601+
[
1602+
(
1603+
'INSERT INTO __ch_gw__.mydb.target ("id") '
1604+
"SELECT __ch_gw__.mydb.source.id FROM __ch_gw__.mydb.source",
1605+
'INSERT INTO "mydb"."target" ("id") SELECT "mydb"."source"."id" FROM "mydb"."source"',
1606+
),
1607+
(
1608+
"SELECT __ch_gw__.mydb.source.id, '__ch_gw__.literal' FROM __ch_gw__.mydb.source "
1609+
"JOIN other_catalog.otherdb.source ON __ch_gw__.mydb.source.id = "
1610+
"other_catalog.otherdb.source.id",
1611+
'SELECT "mydb"."source"."id", \'__ch_gw__.literal\' FROM "mydb"."source" JOIN '
1612+
'"other_catalog"."otherdb"."source" ON "mydb"."source"."id" = '
1613+
'"other_catalog"."otherdb"."source"."id"',
1614+
),
1615+
],
1616+
)
1617+
def test_virtual_catalog_stripped_from_execute_queries(
1618+
make_mocked_engine_adapter: t.Callable, query_sql: str, expected_sql: str
1619+
):
1620+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1621+
adapter.inject_virtual_catalog("ch_gw")
1622+
query = parse_one(query_sql, dialect="clickhouse")
1623+
original_sql = query.sql(dialect="clickhouse")
1624+
1625+
adapter.execute(query)
1626+
1627+
assert query.sql(dialect="clickhouse") == original_sql
1628+
assert to_sql_calls(adapter) == [expected_sql]
1629+
1630+
1631+
def test_virtual_catalog_stripped_from_ctas_and_delete(make_mocked_engine_adapter: t.Callable):
1632+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1633+
adapter.inject_virtual_catalog("ch_gw")
1634+
1635+
adapter.ctas(
1636+
"__ch_gw__.mydb.target",
1637+
parse_one("SELECT __ch_gw__.mydb.source.id FROM __ch_gw__.mydb.source"),
1638+
{"id": exp.DataType.build("Int32")},
1639+
)
1640+
adapter.delete_from(
1641+
"__ch_gw__.mydb.target",
1642+
"__ch_gw__.mydb.target.id IN (SELECT id FROM __ch_gw__.mydb.source)",
1643+
)
1644+
1645+
assert to_sql_calls(adapter) == [
1646+
'CREATE TABLE IF NOT EXISTS "mydb"."target" ENGINE=MergeTree ORDER BY () AS '
1647+
'SELECT CAST("id" AS Nullable(Int32)) AS "id" FROM '
1648+
'(SELECT "mydb"."source"."id" FROM "mydb"."source") AS "_subquery"',
1649+
'DELETE FROM "mydb"."target" WHERE "mydb"."target"."id" IN '
1650+
'(SELECT "id" FROM "mydb"."source")',
1651+
]
1652+
1653+
1654+
def test_virtual_catalog_stripped_from_insert_overwrite(
1655+
make_mocked_engine_adapter: t.Callable, mocker: MockerFixture
1656+
):
1657+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1658+
adapter.inject_virtual_catalog("ch_gw")
1659+
mocker.patch(
1660+
"sqlmesh.core.engine_adapter.EngineAdapter._get_temp_table",
1661+
return_value=exp.to_table("__ch_gw__.mydb.__temp_target_abcd"),
1662+
)
1663+
mocker.patch("sqlmesh.core.engine_adapter.ClickhouseEngineAdapter.fetchone", return_value=None)
1664+
1665+
source_queries, columns_to_types = adapter._get_source_queries_and_columns_to_types(
1666+
parse_one("SELECT * FROM __ch_gw__.mydb.source"),
1667+
{"id": exp.DataType.build("Int8", dialect="clickhouse")},
1668+
"__ch_gw__.mydb.target",
1669+
)
1670+
adapter._insert_overwrite_by_condition(
1671+
"__ch_gw__.mydb.target", source_queries, columns_to_types
1672+
)
1673+
1674+
assert [call.args[0] for call in adapter.cursor.execute.call_args_list] == [
1675+
'CREATE TABLE "mydb"."__temp_target_abcd" AS "mydb"."target"',
1676+
'INSERT INTO "mydb"."__temp_target_abcd" ("id") SELECT "id" FROM '
1677+
'(SELECT * FROM "mydb"."source") AS "_subquery"',
1678+
'EXCHANGE TABLES "mydb"."target" AND "mydb"."__temp_target_abcd"',
1679+
'DROP TABLE IF EXISTS "mydb"."__temp_target_abcd"',
1680+
]
1681+
1682+
1683+
def test_virtual_catalog_stripped_from_rename_table(make_mocked_engine_adapter: t.Callable):
1684+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1685+
adapter.inject_virtual_catalog("ch_gw")
1686+
1687+
adapter.rename_table("__ch_gw__.mydb.old_table", "__ch_gw__.mydb.new_table")
1688+
1689+
assert [call.args[0] for call in adapter.cursor.execute.call_args_list] == [
1690+
'RENAME TABLE "mydb"."old_table" TO "mydb"."new_table"',
1691+
]
1692+
1693+
1694+
def test_virtual_catalog_stripped_from_comments(make_mocked_engine_adapter: t.Callable):
1695+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1696+
adapter.inject_virtual_catalog("ch_gw")
1697+
1698+
adapter._create_table_comment("__ch_gw__.mydb.target", "table comment")
1699+
adapter._create_column_comments("__ch_gw__.mydb.target", {"id": "column comment"})
1700+
1701+
assert [call.args[0] for call in adapter.cursor.execute.call_args_list] == [
1702+
'ALTER TABLE "mydb"."target" MODIFY COMMENT \'table comment\'',
1703+
'ALTER TABLE "mydb"."target" COMMENT COLUMN "id" \'column comment\'',
1704+
]
1705+
1706+
1707+
def test_three_part_names_unchanged_without_virtual_catalog(
1708+
make_mocked_engine_adapter: t.Callable,
1709+
):
1710+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1711+
1712+
adapter.execute(parse_one("SELECT * FROM __ch_gw__.mydb.source", dialect="clickhouse"))
1713+
1714+
assert to_sql_calls(adapter) == ['SELECT * FROM "__ch_gw__"."mydb"."source"']
1715+
1716+
15991717
def test_virtual_catalog_stripped_from_create_view_source(
16001718
make_mocked_engine_adapter: t.Callable,
16011719
):

0 commit comments

Comments
 (0)