Skip to content
Open
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
24 changes: 13 additions & 11 deletions sqlmesh/core/engine_adapter/clickhouse.py
Original file line number Diff line number Diff line change
Expand Up @@ -668,17 +668,19 @@ def _exchange_tables(
f"EXCHANGE TABLES {old_table_sql} AND {new_table_sql}{self._on_cluster_sql()}"
)
except DatabaseError as e:
if "NOT_IMPLEMENTED" in str(e):
# If someone is using an old Clickhouse version, an OS that doesn't support atomic exchanges,
# or a database engine that doesn't support atomic exchanges, we do a non-atomic rename instead.
#
# Executing multiple renames in one call like `RENAME TABLE a to b, c to a` is supported
# but not an atomic operation. Because it is not atomic, doing it in two calls is equivalent
# and does not require defining an additional method.
throwaway_table_name = self._get_temp_table(old_table_name)
self._rename_table(old_table_name, throwaway_table_name)
self._rename_table(new_table_name, old_table_name)
self.drop_table(throwaway_table_name)
if "NOT_IMPLEMENTED" not in str(e):
raise

# If someone is using an old Clickhouse version, an OS that doesn't support atomic exchanges,
# or a database engine that doesn't support atomic exchanges, we do a non-atomic rename instead.
#
# Executing multiple renames in one call like `RENAME TABLE a to b, c to a` is supported
# but not an atomic operation. Because it is not atomic, doing it in two calls is equivalent
# and does not require defining an additional method.
throwaway_table_name = self._get_temp_table(old_table_name)
self._rename_table(old_table_name, throwaway_table_name)
self._rename_table(new_table_name, old_table_name)
self.drop_table(throwaway_table_name)

def _rename_table(
self,
Expand Down
81 changes: 81 additions & 0 deletions tests/core/engine_adapter/test_clickhouse.py
Original file line number Diff line number Diff line change
Expand Up @@ -1428,6 +1428,87 @@ def test_exchange_tables(
]


def _executed_sql(execute_mock: t.Any) -> t.List[str]:
return [
quote_identifiers(call.args[0]).sql("clickhouse")
if isinstance(call.args[0], exp.Expr)
else call.args[0]
for call in execute_mock.call_args_list
]


def test_exchange_tables_reraises_other_errors(
make_mocked_engine_adapter: t.Callable, mocker: MockerFixture, make_temp_table_name: t.Callable
):
from clickhouse_connect.driver.exceptions import DatabaseError # type: ignore

adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)

temp_table_mock = mocker.patch("sqlmesh.core.engine_adapter.EngineAdapter._get_temp_table")
temp_table_mock.return_value = make_temp_table_name("table1", "abcd")

execute_mock = mocker.patch("sqlmesh.core.engine_adapter.ClickhouseEngineAdapter.execute")
execute_mock.side_effect = [
DatabaseError("DB::Exception: Not enough privileges. (ACCESS_DENIED)"),
None,
None,
None,
]

with pytest.raises(DatabaseError, match="ACCESS_DENIED"):
adapter._exchange_tables("table1", "table2")

# No RENAME fallback and no throwaway table drop
assert _executed_sql(execute_mock) == ['EXCHANGE TABLES "table1" AND "table2"']


def test_insert_overwrite_by_condition_replace_exchange_error_propagates(
make_mocked_engine_adapter: t.Callable, mocker: MockerFixture, make_temp_table_name: t.Callable
):
from clickhouse_connect.driver.exceptions import DatabaseError # type: ignore

adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter)

temp_table_mock = mocker.patch("sqlmesh.core.engine_adapter.EngineAdapter._get_temp_table")
temp_table_mock.return_value = make_temp_table_name("target", "abcd")

fetchone_mock = mocker.patch("sqlmesh.core.engine_adapter.ClickhouseEngineAdapter.fetchone")
fetchone_mock.return_value = None

insert_table_name = make_temp_table_name("new_records", "abcd")
existing_table_name = make_temp_table_name("existing_records", "abcd")

source_queries, columns_to_types = adapter._get_source_queries_and_columns_to_types(
parse_one(f"SELECT * FROM {insert_table_name}"),
{
"id": exp.DataType.build("Int8", dialect="clickhouse"),
"ds": exp.DataType.build("Date", dialect="clickhouse"),
},
existing_table_name,
)

def execute_side_effect(sql: t.Any, *args: t.Any, **kwargs: t.Any) -> None:
if str(sql).startswith("EXCHANGE TABLES"):
raise DatabaseError("DB::Exception: Not enough privileges. (ACCESS_DENIED)")

execute_mock = mocker.patch("sqlmesh.core.engine_adapter.ClickhouseEngineAdapter.execute")
execute_mock.side_effect = execute_side_effect

with pytest.raises(DatabaseError, match="ACCESS_DENIED"):
adapter._insert_overwrite_by_condition(
existing_table_name.sql(),
source_queries,
columns_to_types,
)

executed = _executed_sql(execute_mock)
assert executed[-2:] == [
'EXCHANGE TABLES "__temp_existing_records_abcd" AND "__temp_target_abcd"',
'DROP TABLE IF EXISTS "__temp_target_abcd"',
]
assert not any(sql.startswith("RENAME") for sql in executed)


def test_virtual_catalog_ddl_stripping(make_mocked_engine_adapter: t.Callable):
"""After inject_virtual_catalog(), create_schema() with the virtual catalog prefix must strip
the catalog and execute without raising, and with a wrong catalog must raise SQLMeshError."""
Expand Down