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
20 changes: 10 additions & 10 deletions sqlmesh/core/engine_adapter/duckdb.py
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,12 @@ def _create_table(
track_rows_processed: bool = True,
**kwargs: t.Any,
) -> None:
catalog = self.get_current_catalog()
table = (
table_name_or_schema.this
if isinstance(table_name_or_schema, exp.Schema)
else exp.to_table(table_name_or_schema)
)
catalog = table.catalog or self.get_current_catalog()

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

👍 Resolving the catalog from the target table is the right fix here. Small nit: the base adapter already has get_catalog_type_from_table doing the same catalog resolution. If the fix moves into replace_query, following the Trino adapter's pattern, it may be worth reusing that logic or following its shape.

catalog_type_tuple = self.fetchone(
exp.select("type")
.from_("duckdb_databases()")
Expand All @@ -184,6 +189,9 @@ def _create_table(
partitioned_by_exps = None
if catalog_type == "ducklake":
partitioned_by_exps = kwargs.pop("partitioned_by", None)
elif catalog_type == "postgres" and replace:
self.execute(exp.Drop(this=table, kind="TABLE", exists=True, cascade=True))

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This CASCADE will also drop every SQLMesh virtual-layer view that points at this physical table: prod and any dev environments sharing the snapshot. Those views only get recreated during plan promotion, not during sqlmesh run. So after a scheduled run of a FULL model, prod. would disappear until someone applies a plan that re-promotes this snapshot. It would also drop any non-SQLMesh views on this table in Postgres.

DuckDB doesn't run this in a transaction either (SUPPORTS_TRANSACTIONS = False), so if the CTAS below fails, the table's data is lost.

Could we route Postgres-attached targets to the insert-overwrite path instead (see the summary comment), so the table is never dropped?

replace = False

super()._create_table(
table_name_or_schema,
Expand All @@ -199,16 +207,8 @@ def _create_table(
)

if partitioned_by_exps:
# Schema object contains column definitions, so we extract Table
table_name = (
table_name_or_schema.this
if isinstance(table_name_or_schema, exp.Schema)
else table_name_or_schema
)
table_name_str = (
table_name.sql(dialect=self.dialect)
if isinstance(table_name, exp.Table)
else table_name
table.sql(dialect=self.dialect) if isinstance(table, exp.Table) else table

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: table is now always an exp.Table (it comes from exp.to_table or Schema.this), so this isinstance check is dead code. It can just be table.sql(dialect=self.dialect).

)
partitioned_by_str = ", ".join(
expr.sql(dialect=self.dialect) for expr in partitioned_by_exps
Expand Down
22 changes: 22 additions & 0 deletions tests/core/engine_adapter/test_duckdb.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,28 @@ def test_replace_query_pandas(adapter: EngineAdapter, duck_conn):
pd.testing.assert_frame_equal(adapter.fetchdf("SELECT * FROM test_table"), df)


def test_replace_query_attached_postgres(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This checks the generated SQL well. Because fetchone is mocked, though, it can't show how the Postgres extension behaves (whether CASCADE is pushed down, what happens to dependent views). Could we add an integration test with the Docker Postgres engine? For example: attach Postgres, plan a FULL model, run it again, then check that the prod virtual-layer view still exists and can be queried.

make_mocked_engine_adapter: t.Callable, mocker: MockerFixture
) -> None:
adapter = make_mocked_engine_adapter(DuckDBEngineAdapter)
fetchone = mocker.patch.object(adapter, "fetchone", return_value=("postgres",))

adapter.replace_query(
"attached_postgres.test_schema.test_table",
parse_one("SELECT 1 AS a"),
)

assert fetchone.call_count == 1
assert (
fetchone.call_args.args[0].sql(dialect=adapter.dialect)
== "SELECT type FROM DUCKDB_DATABASES() WHERE database_name = 'attached_postgres'"
)
assert to_sql_calls(adapter) == [
'DROP TABLE IF EXISTS "attached_postgres"."test_schema"."test_table" CASCADE',
'CREATE TABLE IF NOT EXISTS "attached_postgres"."test_schema"."test_table" AS SELECT 1 AS "a"',
]


def test_set_current_catalog(make_mocked_engine_adapter: t.Callable, duck_conn):
adapter = make_mocked_engine_adapter(DuckDBEngineAdapter)
adapter.set_current_catalog("test_catalog")
Expand Down
Loading