Loading data
Reads are the goal, so spend effort at load time. A query’s cost is the number of files and row groups it cannot rule out from their min/max statistics, times what it must decode in the rest.
- Set the lake options before the first insert (they apply to new files):
CALL lake.set_option('parquet_row_group_size', 8192); -- see the table below CALL lake.set_option('parquet_compression', 'lz4'); -- decodes faster than snappy - Declare a sort order that matches your filters, before inserting.
DuckLake then sorts every insert and compaction by it:
Put the column queries filter on most first, then the next. For geometry, insert in space-filling-curve order,ALTER TABLE lake.hits SET SORTED BY (CounterID, EventDate, UserID, EventTime);ORDER BY ST_Hilbert(geometry, <extent>), so an area of interest touches few row groups. - Load big tables in key-range batches. One
INSERTof 100M wide rows sorts all of them at once, and that can exhaust memory. Instead, insert disjoint ranges of the leading sort key in order, e.g.WHERE CounterID >= lo AND CounterID < hifor about 5M rows at a time. Each sort stays small, and the ranges don’t overlap across files, so a filtered query opens one or two files instead of all of them. Bound the loader too:SET memory_limit = '4GiB'; -- sorts beyond this spill to temp_directory SET preserve_insertion_order = false; SET temp_directory = '/fast/disk/duckdb-temp'; - Keep columns typed. Convert dates, timestamps and numbers from strings
or epochs at load time (
make_date,epoch_ms, casts). Statistics and comparisons on typed columns prune; on strings they mostly don’t. - Prefer few big inserts to many small ones. Each insert writes its own
files and a catalog snapshot. For trickles, let DuckLake inline small
inserts (
data_inlining_row_limit), then periodically runCALL ducklake_merge_adjacent_files('lake')andducklake_flush_inlined_data. - Partition only on low-cardinality columns (a month, a region) that nearly every query filters on. Each partition value gets its own files, so partitioning on a high-cardinality key makes many tiny files, which is worse than a sort.
Measured on ClickBench hits (100M rows, 105 columns), 13 readers each running
a mix of per-site dashboard queries, PostgreSQL on 3 cores
(scripts/lake.sh hits):
| Load | Load time | Files a site is in | Cold | Warm | New parameters |
|---|---|---|---|---|---|
| Source order, 8K row groups | 62 s | ~all 27 | 16 q/s, p50 337 ms | 40 q/s, p50 140 ms | 25 q/s |
| Sorted in key-range batches, 8K | 139 s | 1 (max 4) | 30 q/s, p50 163 ms | 76 q/s, p50 46 ms | 36 q/s |
| Sorted in key-range batches, 64K | 213 s | 1 (max 4) | 40 q/s, p50 134 ms | 89 q/s, p50 47 ms | 47 q/s |
Sorting doubled read throughput for about twice the load time, and the load never used more than about 4 GiB. Row-group size is a trade-off: on this large sorted table each query reads a contiguous range, so bigger row groups (64K) cut per-row-group overhead and won on throughput. For small lookups of about 1,000 rows (the Houston and lookup benchmarks), 8K was 4× faster than the 122K default and beat 2K. Use 8K for lookup-heavy tables and 32K–64K for large tables queried by ranges of the sort key; the option can be set per table.