Skip to content

DuckDB and the changing physics of analytics

Summary

A guest post by AWS Distinguished Engineer Andy Warfield (intro by Werner Vogels) arguing that the relative costs of compute, memory, and network on a single machine have shifted so far that a large fraction of data work no longer needs to leave the application process. Warfield frames systems design as the art of finding elegant solutions against "changing physics" — the moving relative costs of hardware — and traces how each era's constraints shaped its data systems: MapReduce and Spark were born when single-host I/O bandwidth was small and datasets were large, so they scaled out; today a single m8g.48xlarge has ~50× the memory, cores, and network bandwidth of the beefiest 2007 EC2 instance, which revived interest in extremely efficient single-node engines. DuckDB (created by Hannes Mühleisen and Mark Raasveldt at CWI, the same lab behind MonetDB/X100) embodies this shift: an analytics engine that runs as an in-process library — like SQLite — in the application's own address space. The post announces that DuckLabs, the team behind DuckDB, is joining AWS as a subsidiary, with DuckDB remaining open source (MIT) under the DuckDB Foundation. It also positions DuckDB as complementary to S3's tabular/vector/file primitives (S3 Tables, S3 Vectors, S3 Files), with AWS having sponsored DuckDB's Iceberg extension.

Key takeaways

  1. Systems design chases "changing physics." Unlike the physical sciences (which study immovable truths), computer systems is about elegant compromises against relative hardware costs that keep shifting — memory vs. network speed, compute vs. software-abstraction richness. Success at one scale pushes you into a different set of constraints. (Source: body, "British pub" section)

  2. Distributed query engines were a product of their physics. MapReduce and Spark's RDDs were conceived c. early-2000s under two invariants: datasets were large and the single-host NIC/disk was comparatively small. Their design bet: pay up-front planning / task-shipping / distribution overhead to buy arbitrarily high parallelism and throughput — achieving throughput by adding computers rather than being lean. (Source: body, "limited I/O bandwidth" section)

  3. The hardware physics inverted. A 2007 m1.xlarge (then the beefiest EC2 instance) had 15 GB RAM, 4 vCPU, ~1 Gb/s network. A single m8g.48xlarge today has ~50× more memory, ~50× more cores, and ~50× more network bandwidth. Modern single servers exceed the clusters many first ran Hadoop/Spark on. Even the author's MacBook Pro has 3–5× the cores/RAM, ~40× the memory bandwidth, and >100× the I/O bandwidth of that m1.xlarge. (Source: body, "second computer" section)

  4. Dataset growth is a distribution, and the tail is not the median. The largest datasets grew exponentially, but those are the tail. Most datasets scale with human things — business size, customer count, transactions/day — so the hardware advantage relative to typical data size has grown, reviving efficient single-host engines. (Source: body)

  5. "Scalability! But at what COST?" is the intellectual anchor. McSherry/Isard/Murray (2015) showed a well-optimized single-thread implementation could beat distributed graph systems on 128 cores, and it took 512 cores before the distributed version pulled ahead — a call to pay attention to per-core efficiency. Paul Barham's epigraph: "You can have a second computer once you've shown you know how to use the first one." (Source: body, COST-paper section)

  6. DuckDB's lineage is vectorized single-CPU efficiency. Mühleisen and Raasveldt came from CWI Amsterdam — the MonetDB / X100 lab that rebuilt query execution around cache-resident batches of values once the query-processing bottleneck moved off disk and into the CPU (cache behaviour, branch prediction, superscalar pipelines). DuckDB (started 2018) applied that efficiency focus to an embedded engine; the 2019 SIGMOD demo paper argued for an efficient analytics library packaged like SQLite. (Source: body, "walks like a duck" section)

  7. Library form factor relocates where "work" happens. DuckDB runs in the same address space and works in-situ on the in-memory data structures the app already has, caring as much about its own overhead as about the queries it runs. Library ≠ always-in-the-client: it's a building block placeable at the right point(s) in the stack. The extreme case compiles to WebAssembly and runs entirely in a browser tab (shell.duckdb.org). (Source: body)

  8. Distributed engines are not displaced. "When a job genuinely needs a thousand machines, it needs a thousand machines." The point is that an enormous amount of everyday data work never needed a cluster — and now doesn't have one. Analytics becomes something you do continuously as you build rather than a separate activity elsewhere. (Source: body)

  9. "The glibc of structured data." Warfield's framing: DuckDB is a lean, ubiquitous dependency a great deal of software links against and almost nobody has to think about. Extensibility + format breadth (CSV, Parquet, Delta, Iceberg) + form-factor breadth (browser → server-side engine) makes one efficient codebase raise all ships. (Source: body)

  10. AWS sponsored DuckDB's Iceberg extension; DuckLabs now joins AWS. While building S3 Tables (Iceberg-backed tabular storage), the S3 team saw customers embedding DuckDB directly against their data. AWS became DuckDB's customer/sponsor for the Iceberg extension to broaden Iceberg support beyond the Spark landscape; the extension now implements Iceberg v2 + v3 and motivated the async-I/O support arriving in the 2.0 release. DuckLabs joins AWS as a subsidiary; DuckDB stays open source (MIT) under the DuckDB Foundation, team stays in Amsterdam. (Source: body, "must be a duck" section)

Operational numbers

  • 2007 m1.xlarge: 15 GB RAM, 4 vCPU, ~1 Gb/s network.
  • Today m8g.48xlarge: ~50× memory, ~50× cores, ~50× network vs. m1.xlarge.
  • Author's MacBook Pro vs. m1.xlarge: 3–5× cores/RAM, ~40× memory bandwidth, >100× I/O bandwidth.
  • COST paper (2015): single thread beat 128-core distributed graph systems; 512 cores needed before distributed pulled ahead.
  • DuckDB started 2018; SIGMOD demo paper 2019.
  • DuckDB Iceberg extension: >800K downloads/week; implements Iceberg v2
  • v3; async I/O targets saturating the NIC while scanning S3, arriving in the 2.0 release.
  • DuckDB has maintained development velocity over ~8 years.

Caveats

  • This is partly an acquisition-announcement / advocacy post (DuckLabs joining AWS), not a neutral third-party analysis. The architectural argument stands on its own, but the DuckDB framing is AWS-favorable.
  • No new hard benchmarks of DuckDB itself are presented; the quantitative claims are the hardware-ratio and COST-paper figures.
  • Warfield is explicit that distributed engines remain necessary for genuinely large jobs — the post is a "single-node has grown up," not a "distributed is dead," argument.
  • "~50×" hardware figures are approximate, author-stated round numbers comparing generational EC2 instances.

Source

Last updated · 766 distilled / 2,225 read