DuckDB is usually introduced as "SQLite for analytics", which is true about its packaging and misleading about everything else. It ships as a library you link into your process, with no server, and it stores a database in one file. Inside, it is a columnar, vectorised, parallel query engine with a cost-based optimizer, MVCC transactions and operators that spill to disk when memory runs out. Those internals decide when it is fast, when it runs out of memory, why a second process cannot open your file for writing, and how to size a machine for it.
This article explains the engine from first principles: the path of a query, the on-disk layout, transactions and memory, then a traced query, failure modes and a checklist. For use cases such as querying Parquet in place, read DuckDB practical uses. Version note: as of 2026-10-03 the stable line is 1.5.x and 1.4 is the long-term-support line. Version 2.0 is in alpha, projected for late October 2026; where it changes something below, it is marked as announced, not shipped.
What in-process actually means
A client-server database such as PostgreSQL sends rows to your application over a socket. DuckDB has no such boundary: import duckdb loads the engine into your process, queries run on threads inside it, and results land in your memory. Three consequences follow.
First, there is no transfer cost. DuckDB scans pandas, Polars or Arrow data in place and can return Arrow without per-row conversion. Second, the database lives and dies with your process: committed transactions survive through the write-ahead log, but there is no separate service to monitor or restart. Third, sharing is limited by design. One process at a time may open a database file read-write, and that process may run many writing threads. Several processes may open the same file with access_mode = 'READ_ONLY'. In the 1.5 line, writes across processes go through the Quack remote protocol, which the documentation marks as beta. For production multi-writer lakes it points to DuckLake, which keeps its catalog in a database such as PostgreSQL. For DuckDB 2.0, the project has announced that Quack server mode and a CONNECT statement will graduate to stable. Until you are on a release that ships them, plan on one writer process.
import duckdb
con = duckdb.connect("analytics.duckdb", config={
"memory_limit": "12GB", # default is 80% of RAM: leave room for your own process
"threads": 8, # default is the number of cores
"temp_directory": "/fast-nvme/duckdb_tmp",
"preserve_insertion_order": False, # lets large exports and CTAS stream with less memory
})
# Every connection object shares one database instance; give each thread its own cursor.
def worker(rows):
cur = con.cursor()
cur.executemany("INSERT INTO events VALUES (?, ?, ?)", rows)
The life of a query
The parser turns SQL text into a parse tree. Through the 1.x line it is derived from PostgreSQL's grammar, which is why much PostgreSQL syntax just works; 2.0 announces a new PEG-based parser. The binder resolves names against the catalog, assigns types and expands conveniences such as GROUP BY ALL. The logical planner produces relational operators, and the optimizer rewrites them. It pushes filters and projections down into scans, so a Parquet reader or table scan reads only the needed columns and skips row groups whose statistics rule them out. It decorrelates subqueries, and it orders joins using cardinality estimates.
The physical planner chooses algorithms and cuts the plan into pipelines: chains of operators that stream data, each ending at a pipeline breaker that must see all its input first, such as a hash join build, a hash aggregate, a sort or a window. In the diagram, the dimension table is built into a hash table, the fact table streams through a filter and probes it, and aggregation and sorting finish last.
Vectors: the unit of work
Row-at-a-time engines pay function-call and branch overhead for every value. DuckDB takes the vectorised path pioneered by MonetDB/X100: each operator processes a vector of up to STANDARD_VECTOR_SIZE values, 2,048 by default, from one column. Each interpreted call is amortised over 2,048 values, the compiler can turn the tight loop over a typed array into SIMD code, and the vector stays in L1/L2 cache for the next operator.
A vector has four physical formats: flat (a contiguous array), constant (one value for the whole vector), dictionary (a child vector plus a selection vector of indexes) and sequence (an offset and an increment, as for row ids). Operators work on the compact forms directly, and a filter does not copy surviving rows: it emits a selection vector naming them, which the next operator reads through.
Push-based pipelines and morsel parallelism
Pipelines are driven push-style: a source produces a vector and pushes it through each operator to the sink. Parallelism comes from splitting the source into morsels, chunks such as a range of row groups or a Parquet row group. A pool of worker threads, one per core by default, takes morsels as each finishes the last one. Morsels are small and handed out dynamically, so a slow core or a skewed morsel does not stall the others.
Sinks need care: a parallel hash aggregate gives each thread a local, hash-partitioned table and merges partitions in parallel, and a sort merges per-thread runs. This is why DuckDB scales with cores and why threads is the first knob to check. It also explains two surprises: preserve_insertion_order (default true) costs buffering on large exports, and a LIMIT without ORDER BY can return different rows each run.
Storage: one file, row groups and compressed segments
A DuckDB database is a single file of fixed-size blocks, 262,144 bytes (256 KiB) by default. Tables are cut horizontally into row groups of up to 122,880 rows, sixty vectors, and each row group stores each column as its own compressed segment. Each segment carries statistics including minimum and maximum, which act as a zonemap. A scan with a selective predicate on a column that correlates with load order, such as a timestamp, skips most segments without decompressing them. That is the same pruning idea columnar engines in general rely on.
Compression is chosen per segment at checkpoint time by analysing the data: constant, run-length, bit packing, frame of reference, dictionary, FSST for strings, ALP, Chimp and Patas for floats, and Zstd. You do not declare encodings, but you can inspect them:
-- Which compression did each column segment get, and what do its statistics say?
SELECT row_group_id, column_name, segment_type, compression, count, stats
FROM pragma_storage_info('orders')
WHERE column_name IN ('status', 'amount')
ORDER BY row_group_id
LIMIT 8;
-- Fold the WAL into the main file now (it also happens automatically at checkpoint_threshold).
CHECKPOINT;Durability uses a write-ahead log. A commit appends its changes to file.duckdb.wal. When the log passes checkpoint_threshold (16 MiB by default), or when you run CHECKPOINT or close the database cleanly, a checkpoint rewrites the affected row groups into the main file and truncates the WAL.
Since v0.10, newer releases read older files. For v1.0 through v1.5, new files default to storage version 64, the v1.0.0 format, so older 1.x readers can open them unless you opt in with ATTACH ... (STORAGE_VERSION ...). The 2.0 announcement moves the default to v2.0.0; assume 1.x cannot read such files until you have tested it.
Transactions: MVCC with optimistic conflicts
DuckDB provides ACID transactions through multi-version concurrency control modelled on HyPer's: updates happen in place and previous versions go to undo buffers that older transactions can still read. Each transaction sees the database as of its start, so long reads never block writers or the reverse.
Conflicts are optimistic. Appends never conflict, even on one table, and updates to different rows proceed concurrently. If two transactions modify the same row, the second fails with a conflict error and must retry. Partition writers by key and wrap the rest in a retry loop:
import random, time
import duckdb
def run_with_retry(con, sql, params=(), attempts=5):
"""Optimistic concurrency: two writers touching the same rows -> the later one fails."""
for i in range(attempts):
cur = con.cursor()
try:
cur.execute("BEGIN TRANSACTION")
cur.execute(sql, params)
cur.execute("COMMIT")
return
except duckdb.TransactionException:
cur.execute("ROLLBACK")
time.sleep((2 ** i) * 0.01 + random.random() * 0.01)
raise RuntimeError("gave up after repeated write conflicts")
Memory: the buffer manager and spilling
All persistent blocks and large temporary structures go through a buffer manager capped at memory_limit, 80% of physical RAM by default. Operators pin the blocks they need and unpin them afterwards. When the limit is reached, unpinned blocks are evicted: clean data blocks are dropped, and temporary data is written to the temp directory, capped by max_temp_directory_size (90% of free disk by default). Hash aggregates, hash joins, sorts and window functions have out-of-core versions, so a query larger than RAM usually slows down rather than failing.
"Usually" has limits. The memory limit covers the buffer manager, not your whole process: a dataframe you materialise or a fetched result sits outside it, so a container sized exactly to memory_limit can still be killed. Huge list() aggregates and wide pivots spill poorly, and a slow or tiny temp disk turns spilling into a stall or an out-of-space error.
Worked example: reading a profile
Take a 400-million-row orders table loaded in time order, a 2-million-row customers table, and a 16 GB laptop. The question: daily revenue per region for September.
EXPLAIN ANALYZE
SELECT c.region, date_trunc('day', o.created_at) AS day, sum(o.amount) AS revenue
FROM orders o
JOIN customers c USING (customer_id)
WHERE o.created_at >= DATE '2026-09-01'
GROUP BY ALL
ORDER BY revenue DESC;Read the profile bottom-up. The TABLE_SCAN on orders should show the date filter pushed into it and a row count far below 400 million. Because the table was loaded in time order, zonemaps skip almost every row group outside September. If the scan reports all 400 million rows, the data is not clustered on the filter column, and an ORDER BY created_at during the load would fix it.
Next, the HASH_JOIN should build on customers, the smaller side; if it built on orders, a cardinality estimate was wrong. The HASH_GROUP_BY has few groups, so it stays in memory. If most time is in the scan with busy CPUs, more threads help; if CPUs idle, you are waiting on I/O, typically remote Parquet.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
| Could not set lock on file | A second process opened the file read-write | One writer process; readers with READ_ONLY; or a server layer |
| Transaction conflict on update | Two writers modified the same rows | Partition writes by key; retry with backoff |
| Out of memory despite spilling | Memory outside the buffer manager, or a non-spilling structure | Lower memory_limit below the container limit; stream results with fetch batches |
| Query stalls on large join | Spilling to a slow or full temp directory | Point temp_directory at fast local disk; check max_temp_directory_size |
| File grows after deletes | Space is reused, not returned to the OS | Copy into a fresh database when the size matters |
| Old reader cannot open file | File written with a newer storage version | Pin STORAGE_VERSION for shared files; upgrade readers |
| Different rows each run | LIMIT without ORDER BY under parallel execution | Always order when the rows matter |
Trade-offs
Against SQLite, DuckDB trades fast single-row transactional writes for fast scans and aggregates. Use SQLite for an application's operational state and DuckDB for questions over many rows; the split is the one in OLTP versus OLAP. Against a server engine such as ClickHouse, DuckDB has no replication, no cluster, no multi-tenant resource control and, in the 1.x line, no supported multi-process writers. In exchange it has zero operations and zero transfer cost. The honest boundary is concurrency and shared state, not data size: one machine with a fast disk handles hundreds of gigabytes, but not forty analysts writing the same tables.
What to do next
- Check your version with
SELECT version()and decide between the 1.4 LTS line and the current 1.5 line. Test 2.0 only on copies of your files until it ships and you have verified readers. - Set
memory_limitexplicitly, below your container or VM limit, and pointtemp_directoryat fast local disk. - Run
EXPLAIN ANALYZEon your three most expensive queries and confirm that filters reach the scans and that hash joins build on the smaller side. - Load large tables sorted by the column you filter on most, then inspect
pragma_storage_infoto see segment statistics and compression. - Audit how processes open each file: exactly one read-write process, all others
READ_ONLY. - Wrap concurrent updates in a conflict-retry loop and partition writers by key.
- Run
CHECKPOINTafter bulk loads, and pinSTORAGE_VERSIONon files other tools read.