Skip to content

Commit c287055

Browse files
committed
fix(clickhouse): strip virtual catalog from string-built statements
`execute()` sends string statements verbatim, so the central stripping in `_to_sql` never sees them. `_create_table_like`, `_exchange_tables`, `_rename_table` and the table/column comment builders still sent the injected virtual catalog to ClickHouse. This broke insert-overwrite (incremental model loads) and silently dropped comments. Route those table names through `_strip_virtual_catalog`, and parse string names with the ClickHouse dialect when stripping. `_create_table_like` now renders quoted identifiers, consistent with the rename/exchange statements. Also add a regression test that three-part names are left untouched when no virtual catalog has been injected. Signed-off-by: mday-io <mdaytn@gmail.com>
1 parent 9ceed6e commit c287055

2 files changed

Lines changed: 93 additions & 16 deletions

File tree

‎sqlmesh/core/engine_adapter/clickhouse.py‎

Lines changed: 22 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -517,8 +517,14 @@ def _create_table_like(
517517
**kwargs: t.Any,
518518
) -> None:
519519
"""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+
)
520526
self.execute(
521-
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}"
522528
)
523529

524530
def _get_partition_ids(
@@ -663,7 +669,7 @@ def _strip_virtual_catalog(self, name: "TableName") -> exp.Table:
663669
SQL is sent to the wire, since ClickHouse only supports a two-level
664670
``[database].[table]`` naming scheme.
665671
"""
666-
table = exp.to_table(name)
672+
table = exp.to_table(name, dialect=self.dialect)
667673
if self._default_catalog and table.catalog == self._default_catalog:
668674
table.set("catalog", None)
669675
return table
@@ -675,8 +681,12 @@ def _exchange_tables(
675681
) -> None:
676682
from clickhouse_connect.driver.exceptions import DatabaseError # type: ignore
677683

678-
old_table_sql = exp.to_table(old_table_name).sql(dialect=self.dialect, identify=True)
679-
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+
)
680690

681691
try:
682692
self.execute(
@@ -700,8 +710,12 @@ def _rename_table(
700710
old_table_name: TableName,
701711
new_table_name: TableName,
702712
) -> None:
703-
old_table_sql = exp.to_table(old_table_name).sql(dialect=self.dialect, identify=True)
704-
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+
)
705719

706720
self.execute(f"RENAME TABLE {old_table_sql} TO {new_table_sql}{self._on_cluster_sql()}")
707721

@@ -989,7 +1003,7 @@ def _build_view_properties_exp(
9891003
def _build_create_comment_table_exp(
9901004
self, table: exp.Table, table_comment: str, table_kind: str, **kwargs: t.Any
9911005
) -> exp.Comment | str:
992-
table_sql = table.sql(dialect=self.dialect, identify=True)
1006+
table_sql = self._strip_virtual_catalog(table).sql(dialect=self.dialect, identify=True)
9931007

9941008
truncated_comment = self._truncate_table_comment(table_comment)
9951009
comment_sql = exp.Literal.string(truncated_comment).sql(dialect=self.dialect)
@@ -1004,7 +1018,7 @@ def _build_create_comment_column_exp(
10041018
table_kind: str = "TABLE",
10051019
**kwargs: t.Any,
10061020
) -> exp.Comment | str:
1007-
table_sql = table.sql(dialect=self.dialect, identify=True)
1021+
table_sql = self._strip_virtual_catalog(table).sql(dialect=self.dialect, identify=True)
10081022
column_sql = exp.to_column(column_name).sql(dialect=self.dialect, identify=True)
10091023

10101024
truncated_comment = self._truncate_table_comment(column_comment)

‎tests/core/engine_adapter/test_clickhouse.py‎

Lines changed: 71 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"',
@@ -1655,6 +1655,69 @@ def test_virtual_catalog_stripped_from_ctas_and_delete(make_mocked_engine_adapte
16551655
]
16561656

16571657

1658+
def test_virtual_catalog_stripped_from_insert_overwrite(
1659+
make_mocked_engine_adapter: t.Callable, mocker: MockerFixture
1660+
):
1661+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1662+
adapter.inject_virtual_catalog("ch_gw")
1663+
mocker.patch(
1664+
"sqlmesh.core.engine_adapter.EngineAdapter._get_temp_table",
1665+
return_value=exp.to_table("__ch_gw__.mydb.__temp_target_abcd"),
1666+
)
1667+
mocker.patch("sqlmesh.core.engine_adapter.ClickhouseEngineAdapter.fetchone", return_value=None)
1668+
1669+
source_queries, columns_to_types = adapter._get_source_queries_and_columns_to_types(
1670+
parse_one("SELECT * FROM __ch_gw__.mydb.source"),
1671+
{"id": exp.DataType.build("Int8", dialect="clickhouse")},
1672+
"__ch_gw__.mydb.target",
1673+
)
1674+
adapter._insert_overwrite_by_condition(
1675+
"__ch_gw__.mydb.target", source_queries, columns_to_types
1676+
)
1677+
1678+
assert [call.args[0] for call in adapter.cursor.execute.call_args_list] == [
1679+
'CREATE TABLE "mydb"."__temp_target_abcd" AS "mydb"."target"',
1680+
'INSERT INTO "mydb"."__temp_target_abcd" ("id") SELECT "id" FROM '
1681+
'(SELECT * FROM "mydb"."source") AS "_subquery"',
1682+
'EXCHANGE TABLES "mydb"."target" AND "mydb"."__temp_target_abcd"',
1683+
'DROP TABLE IF EXISTS "mydb"."__temp_target_abcd"',
1684+
]
1685+
1686+
1687+
def test_virtual_catalog_stripped_from_rename_table(make_mocked_engine_adapter: t.Callable):
1688+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1689+
adapter.inject_virtual_catalog("ch_gw")
1690+
1691+
adapter.rename_table("__ch_gw__.mydb.old_table", "__ch_gw__.mydb.new_table")
1692+
1693+
assert [call.args[0] for call in adapter.cursor.execute.call_args_list] == [
1694+
'RENAME TABLE "mydb"."old_table" TO "mydb"."new_table"',
1695+
]
1696+
1697+
1698+
def test_virtual_catalog_stripped_from_comments(make_mocked_engine_adapter: t.Callable):
1699+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1700+
adapter.inject_virtual_catalog("ch_gw")
1701+
1702+
adapter._create_table_comment("__ch_gw__.mydb.target", "table comment")
1703+
adapter._create_column_comments("__ch_gw__.mydb.target", {"id": "column comment"})
1704+
1705+
assert [call.args[0] for call in adapter.cursor.execute.call_args_list] == [
1706+
'ALTER TABLE "mydb"."target" MODIFY COMMENT \'table comment\'',
1707+
'ALTER TABLE "mydb"."target" COMMENT COLUMN "id" \'column comment\'',
1708+
]
1709+
1710+
1711+
def test_three_part_names_unchanged_without_virtual_catalog(
1712+
make_mocked_engine_adapter: t.Callable,
1713+
):
1714+
adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)
1715+
1716+
adapter.execute(parse_one("SELECT * FROM __ch_gw__.mydb.source", dialect="clickhouse"))
1717+
1718+
assert to_sql_calls(adapter) == ['SELECT * FROM "__ch_gw__"."mydb"."source"']
1719+
1720+
16581721
def test_virtual_catalog_stripped_from_create_view_source(
16591722
make_mocked_engine_adapter: t.Callable,
16601723
):

0 commit comments

Comments
 (0)