Repository navigation
Fix replacement in attached Postgres catalogs #6078
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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() | ||
| catalog_type_tuple = self.fetchone( | ||
| exp.select("type") | ||
| .from_("duckdb_databases()") | ||
|
|
@@ -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)) | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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, | ||
|
|
@@ -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 | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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( | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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") | ||
|
|
||
There was a problem hiding this comment.
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.