Skip to content

Commit 5af58f8

Browse files
committed
fix(duckdb): partition initial DuckLake writes
Signed-off-by: Anas Khan <83116240+anxkhn@users.noreply.github.com>
1 parent d69c262 commit 5af58f8

2 files changed

Lines changed: 57 additions & 7 deletions

File tree

‎sqlmesh/core/engine_adapter/duckdb.py‎

Lines changed: 28 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -173,9 +173,30 @@ def _create_table(
173173
track_rows_processed: bool = True,
174174
**kwargs: t.Any,
175175
) -> None:
176+
table_name = (
177+
table_name_or_schema.this
178+
if isinstance(table_name_or_schema, exp.Schema)
179+
else table_name_or_schema
180+
)
181+
catalog = exp.to_table(table_name).catalog or self.get_current_catalog()
182+
176183
partitioned_by_exps = None
177-
if self._get_catalog_type(self.get_current_catalog()) == "ducklake":
184+
insert_expression = None
185+
if self._get_catalog_type(catalog) == "ducklake":
178186
partitioned_by_exps = kwargs.pop("partitioned_by", None)
187+
if (
188+
partitioned_by_exps
189+
and expression is not None
190+
and (replace or not exists or not self.table_exists(table_name))
191+
):
192+
insert_expression = expression.copy()
193+
query = t.cast(exp.Query, expression)
194+
expression = (
195+
exp.select("*")
196+
.from_(query.subquery("_sqlmesh_schema_only", copy=False))
197+
.where(exp.false())
198+
.limit(0)
199+
)
179200

180201
super()._create_table(
181202
table_name_or_schema,
@@ -191,12 +212,6 @@ def _create_table(
191212
)
192213

193214
if partitioned_by_exps:
194-
# Schema object contains column definitions, so we extract Table
195-
table_name = (
196-
table_name_or_schema.this
197-
if isinstance(table_name_or_schema, exp.Schema)
198-
else table_name_or_schema
199-
)
200215
table_name_str = (
201216
table_name.sql(dialect=self.dialect)
202217
if isinstance(table_name, exp.Table)
@@ -207,6 +222,12 @@ def _create_table(
207222
)
208223
self.execute(f"ALTER TABLE {table_name_str} SET PARTITIONED BY ({partitioned_by_str});")
209224

225+
if insert_expression is not None:
226+
self.execute(
227+
exp.insert(insert_expression, exp.to_table(table_name)),
228+
track_rows_processed=track_rows_processed,
229+
)
230+
210231
def _drop_object(
211232
self,
212233
name: TableName | SchemaName,

‎tests/core/engine_adapter/test_duckdb.py‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import typing as t
2+
from datetime import date
23

34
import pandas as pd # noqa: TID253
45
import pytest
@@ -226,3 +227,31 @@ def test_drop_object_cascade_by_catalog_type(make_mocked_engine_adapter: t.Calla
226227
'DROP SCHEMA IF EXISTS "lake"."virt" CASCADE',
227228
'DROP TABLE IF EXISTS "native"."phys"."t" CASCADE',
228229
]
230+
231+
232+
def test_ducklake_partitioning_on_initial_query(adapter: EngineAdapter, duck_conn, tmp_path):
233+
catalog = "ducklake_initial_partition_db"
234+
235+
duck_conn.install_extension("ducklake")
236+
duck_conn.load_extension("ducklake")
237+
duck_conn.execute(
238+
f"ATTACH 'ducklake:{tmp_path}/{catalog}.ducklake' AS {catalog} "
239+
f"(DATA_PATH '{tmp_path}', DATA_INLINING_ROW_LIMIT 0);"
240+
)
241+
242+
adapter.create_schema(f"{catalog}.test_schema")
243+
adapter.replace_query(
244+
f"{catalog}.test_schema.test_table",
245+
parse_one("SELECT 1 AS id, DATE '2000-01-01' AS ds UNION ALL SELECT 2, DATE '2000-01-02'"),
246+
partitioned_by=[exp.to_column("ds")],
247+
)
248+
249+
assert adapter.fetchall(f"SELECT * FROM {catalog}.test_schema.test_table ORDER BY id") == [
250+
(1, date(2000, 1, 1)),
251+
(2, date(2000, 1, 2)),
252+
]
253+
partition_ids = duck_conn.execute(
254+
f"SELECT partition_id FROM __ducklake_metadata_{catalog}.main.ducklake_data_file"
255+
).fetchall()
256+
assert partition_ids
257+
assert all(partition_id is not None for (partition_id,) in partition_ids)

0 commit comments

Comments
 (0)