Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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.

  1. 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
    
  2. Declare a sort order that matches your filters, before inserting. DuckLake then sorts every insert and compaction by it:
    ALTER TABLE lake.hits SET SORTED BY (CounterID, EventDate, UserID, EventTime);
    
    Put the column queries filter on most first, then the next. For geometry, insert in space-filling-curve order, ORDER BY ST_Hilbert(geometry, <extent>), so an area of interest touches few row groups.
  3. Load big tables in key-range batches. One INSERT of 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 < hi for 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';
    
  4. 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.
  5. 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 run CALL ducklake_merge_adjacent_files('lake') and ducklake_flush_inlined_data.
  6. 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):

LoadLoad timeFiles a site is inColdWarmNew parameters
Source order, 8K row groups62 s~all 2716 q/s, p50 337 ms40 q/s, p50 140 ms25 q/s
Sorted in key-range batches, 8K139 s1 (max 4)30 q/s, p50 163 ms76 q/s, p50 46 ms36 q/s
Sorted in key-range batches, 64K213 s1 (max 4)40 q/s, p50 134 ms89 q/s, p50 47 ms47 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.