Skip to content

Commit dfa099e

Browse files
Fix replacement in attached Postgres catalogs
Signed-off-by: clowrance-proppilot <clowrance@propertypilot.com>
1 parent 76225c6 commit dfa099e

2 files changed

Lines changed: 32 additions & 10 deletions

File tree

sqlmesh/core/engine_adapter/duckdb.py

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,12 @@ def _create_table(
173173
track_rows_processed: bool = True,
174174
**kwargs: t.Any,
175175
) -> None:
176-
catalog = self.get_current_catalog()
176+
table = (
177+
table_name_or_schema.this
178+
if isinstance(table_name_or_schema, exp.Schema)
179+
else exp.to_table(table_name_or_schema)
180+
)
181+
catalog = table.catalog or self.get_current_catalog()
177182
catalog_type_tuple = self.fetchone(
178183
exp.select("type")
179184
.from_("duckdb_databases()")
@@ -184,6 +189,9 @@ def _create_table(
184189
partitioned_by_exps = None
185190
if catalog_type == "ducklake":
186191
partitioned_by_exps = kwargs.pop("partitioned_by", None)
192+
elif catalog_type == "postgres" and replace:
193+
self.execute(exp.Drop(this=table, kind="TABLE", exists=True, cascade=True))
194+
replace = False
187195

188196
super()._create_table(
189197
table_name_or_schema,
@@ -199,16 +207,8 @@ def _create_table(
199207
)
200208

201209
if partitioned_by_exps:
202-
# Schema object contains column definitions, so we extract Table
203-
table_name = (
204-
table_name_or_schema.this
205-
if isinstance(table_name_or_schema, exp.Schema)
206-
else table_name_or_schema
207-
)
208210
table_name_str = (
209-
table_name.sql(dialect=self.dialect)
210-
if isinstance(table_name, exp.Table)
211-
else table_name
211+
table.sql(dialect=self.dialect) if isinstance(table, exp.Table) else table
212212
)
213213
partitioned_by_str = ", ".join(
214214
expr.sql(dialect=self.dialect) for expr in partitioned_by_exps

tests/core/engine_adapter/test_duckdb.py

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,28 @@ def test_replace_query_pandas(adapter: EngineAdapter, duck_conn):
6868
pd.testing.assert_frame_equal(adapter.fetchdf("SELECT * FROM test_table"), df)
6969

7070

71+
def test_replace_query_attached_postgres(
72+
make_mocked_engine_adapter: t.Callable, mocker: MockerFixture
73+
) -> None:
74+
adapter = make_mocked_engine_adapter(DuckDBEngineAdapter)
75+
fetchone = mocker.patch.object(adapter, "fetchone", return_value=("postgres",))
76+
77+
adapter.replace_query(
78+
"attached_postgres.test_schema.test_table",
79+
parse_one("SELECT 1 AS a"),
80+
)
81+
82+
assert fetchone.call_count == 1
83+
assert (
84+
fetchone.call_args.args[0].sql(dialect=adapter.dialect)
85+
== "SELECT type FROM DUCKDB_DATABASES() WHERE database_name = 'attached_postgres'"
86+
)
87+
assert to_sql_calls(adapter) == [
88+
'DROP TABLE IF EXISTS "attached_postgres"."test_schema"."test_table" CASCADE',
89+
'CREATE TABLE IF NOT EXISTS "attached_postgres"."test_schema"."test_table" AS SELECT 1 AS "a"',
90+
]
91+
92+
7193
def test_set_current_catalog(make_mocked_engine_adapter: t.Callable, duck_conn):
7294
adapter = make_mocked_engine_adapter(DuckDBEngineAdapter)
7395
adapter.set_current_catalog("test_catalog")

0 commit comments

Comments
 (0)