System Design
B12

Time series and analytics: partition, roll up, scan columns

Metrics and events arrive in time order and are read in time windows. Partition by time, pre-aggregate what dashboards read, and move heavy scans to a column store.

Not startedSaved in this browser only.
  1. 1OLTP reads and writes a few rows by key. OLAP scans one or two columns of millions of rows.
  2. 2A column store reads only the columns a query names. Similar values side by side compress well.
  3. 3Partition time series by day. Recent data is hot; drop old days whole.
  4. 4Dashboards read rollups that store count and sum, not raw rows.
B12
    A

    Rows or columns

    one question, two layouts
    SELECT avg(cpu) · one day · 432,000 samplesRow store8 KB pages, rows side by sidepage 1page 2Reads every column of every row:8.64 MB20 B a row × 432,000Column storeone file per columntshostcpuReads only the cpu column:3.46 MB plain0.81 MB encodeddelta encoding
    • OLTP: many short transactions, a few rows each, by key. Store rows.
    • OLAP: few large queries, millions of rows, few columns. Store columns.
    • Time series: append mostly, read by time window, recent data most.

    A row store reads every column of the rows it scans. A column store reads only the columns the query names, and those compress well.

    B

    Where the dashboard reads

    for a metrics dashboard
    methodwherestatus
    Scan raw rows on the OLTP primaryPostgresNot approved
    Raw rows, partitioned by dayPostgresShort windows
    Rollup tables, upserted by a jobPostgresApproved
    Materialised view, full refreshPostgresShort history
    Continuous aggregatesTimescaleDBApproved
    Column store, fed by CDCClickHouseApproved
    Host metrics with bounded labelsPrometheusApproved
    A 1% sample of pagesPostgresTotals, averages

    A full refresh recomputes every bucket of history. A rollup job writes only the last window.

    I keep raw samples partitioned by day, and the dashboard reads rollups. A column store comes in when ad hoc scans outgrow Postgres.

    C

    Schema

    partitions, then rollups
    metrics, partitioned by day (UTC)10-0210-0310-0410-0510-0610-0710-0810-0910-10last 24 h: 2 partitionsaheaddrophotresolution tiers · rows for 7.5 days, 50 hostsraw, every 10 s3.24 M rows · 259.8 MBkeep 30 daysper minute540 k rows · 62 MBkeep 1 yearper hour9 k rows · 1 MBkeep for yearsbar length on a log scale · retention periods are an example policy
    rows
    50 hosts × 1 sample per 10 s × 7.5 days = 3,240,000
    row size
    52.3 bytes per row in the heap, for 20 bytes of data
    index
    (host, ts) on each partition: 98.3 MB beside 161.5 MB of rows
    raw samples and rollupssql
    CREATE TABLE metrics (
      ts    timestamptz      NOT NULL,
      host  int              NOT NULL,
      value double precision NOT NULL
    ) PARTITION BY RANGE (ts)1;
    
    CREATE INDEX ON metrics (host, ts)2;
    
    CREATE TABLE metrics_1m (
      minute timestamptz      NOT NULL,
      host   int              NOT NULL,
      n      bigint           NOT NULL,
      total  double precision NOT NULL4,
      lo     double precision NOT NULL,
      hi     double precision NOT NULL,
      PRIMARY KEY (minute, host)3
    );
    
    CREATE TABLE metrics_1h (
      hour  timestamptz      NOT NULL,
      host  int              NOT NULL,
      n     bigint           NOT NULL,
      total double precision NOT NULL,
      lo    double precision NOT NULL,
      hi    double precision NOT NULL,
      PRIMARY KEY (hour, host)
    );
    
    CREATE FUNCTION add_day_partition5(day date) RETURNS void LANGUAGE plpgsql AS $$
    BEGIN
      EXECUTE format('CREATE TABLE IF NOT EXISTS %I PARTITION OF metrics FOR VALUES FROM (%L) TO (%L)',
        'metrics_' || to_char(day, 'YYYYMMDD'),
        day::timestamp AT TIME ZONE 'UTC',
        (day + 1)::timestamp AT TIME ZONE 'UTC');
    END $$;
    1. 1One child table per day. A filter on ts skips the other days, and retention drops a day whole.
    2. 2Created on the parent, so every partition gets its own copy. One host over a day is an index range scan.
    3. 3One row per host per minute. The rollup job upserts on this key, so a rerun replaces the row.
    4. 4Store the sum and the count, not the average. Sums add up across buckets; averages do not.
    5. 5A daily job creates tomorrow before any sample for it arrives. A write with no partition fails.

    Raw samples go in one partition per day. Rollups keep count, sum, min and max per minute and per hour, so they combine.

    D

    The warehouse pattern

    click a step; its path lights up
    Serviceorders, paymentsPostgresOLTP primaryCDC readerlogical decodingKafkachange logColumn storeClickHouse, BigQueryDashboardsand analysts

    Step 1: Write

    • The service writes orders in short transactions.
    • Each one touches a few rows by key.
    • No dashboard query runs here.

    If it fails

    Nothing new: the primary keeps its own durability. Analytics never slows a write.

    I keep the OLTP primary for transactions. Changes flow from its WAL through a log into a column store, and analytics runs there.

    E

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    PostgresDeclarative partitioning: PARTITION BY RANGEOne table per day, behind one table name.Multi-tenant splits, archives
    PostgresPartition pruningA filter on ts reads 2 of 9 partitions for the last 24 h.Any range or list partition key
    PostgresDROP or DETACH a partitionRetention in 0.79 ms, with no row scan and no dead rows.Archiving to cheap storage
    Postgresdate_bin(interval, ts, origin)Fixed buckets from a UTC origin, whatever the session time zone.Any time-bucketed report
    PostgresINSERT ... ON CONFLICT DO UPDATEThe rollup job replaces a bucket, so a rerun or a late sample is safe.Idempotent writes, counters
    PostgresMaterialised view, REFRESH CONCURRENTLYA stored query result that reads keep using during a refresh. It needs a unique index.Expensive reports
    PostgresTABLESAMPLE SYSTEM (1)Reads 1% of pages for an approximate total or average.Data exploration
    PostgresLogical decoding, replication slotsA CDC reader gets every committed change from the WAL.Search indexing, cache invalidation
    PostgresRow storageLimit A scan reads every column; 52.3 bytes per sample here.
    TimescaleDBHypertables, continuous aggregates, columnar compressionDay partitions, rollups and compression as one Postgres extension.IoT, finance ticks
    ClickHouseMergeTree: sorted column parts, codecs (Delta, DoubleDelta, Gorilla), TTLCompressed columns, sorted by the key a filter uses, with old rows expired.Logs, product analytics
    ClickHouseMaterialised views that run on insertRollups maintained per batch, with no refresh.Real-time counters
    BigQueryServerless columns, billed by bytes scanned; partitioning and clusteringAd hoc SQL over all history. A partition filter cuts the bytes billed.Warehouse, ML features
    BigQueryAPPROX_COUNT_DISTINCT (HyperLogLog++)Distinct users in a fixed small memory, with a small error.Unique visitors
    PrometheusPull scraping, 2-hour blocks, recording rules, 15-day default retentionHost and service metrics at 1 to 2 bytes per sample.Alerting on SLOs

    Tool facts from each project's documentation: PostgreSQL, TimescaleDB, ClickHouse, BigQuery and Prometheus.

    Postgres gives me partitions, pruning and upserts for rollups. A column store gives me compressed columns and aggregates on insert. Prometheus is for bounded host metrics.

    F

    Roll up, keep, drop

    pseudo code
    the upkeep jobspseudo code
    every minute, for [now - 5 min, now):          // re-read a little, for late samples1
      group raw samples by (minute, host)
      upsert into the minute rollup: count, sum, min, max   // replace the row, never add to it2
    
    every hour, for the last hour:
      group minute rows by (hour, host) into the hour rollup
    
    dashboard: average = sum(total) / sum(n)              // never an average of averages3
    
    every day:
      create the partition for tomorrow                   // a write never finds no partition
      drop the raw partition older than 30 days           // one file removed, no row scan4
    1. 1Samples can arrive minutes late. Each run covers a few closed minutes again.
    2. 2The upsert writes the full count for the bucket. A rerun gives the same row, so retries are safe.
    3. 3Divide the summed totals by the summed counts. Buckets with more samples weigh more.
    4. 40.79 ms in the lab, against 454 ms to DELETE the same day row by row.
    Tested source SQL: minute and hour rollups · SQL: partitions · SQL: dashboard on the rollup
    SQL: minute and hour rollupssql
    INSERT INTO metrics_1m (minute, host, n, total, lo, hi)
    SELECT date_bin('1 minute', ts, '2000-01-01 00:00+00'), host, count(*), sum(value), min(value), max(value)
    FROM metrics
    WHERE ts >= $1 AND ts < $2
    GROUP BY 1, 2
    ON CONFLICT (minute, host) DO UPDATE
    SET n = EXCLUDED.n, total = EXCLUDED.total, lo = EXCLUDED.lo, hi = EXCLUDED.hi;
    
    INSERT INTO metrics_1h (hour, host, n, total, lo, hi)
    SELECT date_bin('1 hour', minute, '2000-01-01 00:00+00'), host, sum(n), sum(total), min(lo), max(hi)
    FROM metrics_1m
    WHERE minute >= $1 AND minute < $2
    GROUP BY 1, 2
    ON CONFLICT (hour, host) DO UPDATE
    SET n = EXCLUDED.n, total = EXCLUDED.total, lo = EXCLUDED.lo, hi = EXCLUDED.hi;
    SQL: partitionssql
    CREATE FUNCTION add_day_partition(day date) RETURNS void LANGUAGE plpgsql AS $$
    BEGIN
      EXECUTE format('CREATE TABLE IF NOT EXISTS %I PARTITION OF metrics FOR VALUES FROM (%L) TO (%L)',
        'metrics_' || to_char(day, 'YYYYMMDD'),
        day::timestamp AT TIME ZONE 'UTC',
        (day + 1)::timestamp AT TIME ZONE 'UTC');
    END $$;
    
    CREATE FUNCTION drop_day_partition(day date) RETURNS void LANGUAGE plpgsql AS $$
    BEGIN
      EXECUTE format('DROP TABLE IF EXISTS %I', 'metrics_' || to_char(day, 'YYYYMMDD'));
    END $$;
    SQL: dashboard on the rollupsql
    SELECT minute AS bucket, sum(n) AS n, sum(total) AS total, min(lo) AS lo, max(hi) AS hi
    FROM metrics_1m
    WHERE minute >= $1 AND minute < $2
    GROUP BY 1 ORDER BY 1;

    A job upserts count, sum, min and max per minute, then per hour. A daily job creates tomorrow's partition and drops the oldest.

    G

    Try it: query cost

    recorded from Postgres
    query
    read from
    median time
    147 ms
    rows read
    648,000
    pages read
    32.3 MB
    partitions
    2 of 9
    answer
    1,440 points, equal to the raw rows
    rows read, log scale1k10k100k1M10MRaw rows147 msMinute rollup29 msMaterialised view25 msHour rollupno minute buckets

    Postgres 16, one process per query, warm cache. Median of 5 runs. The view refresh takes 5,692 ms; rolling up one hour takes 56 ms.

    On the minute rollup the 24-hour dashboard reads 9 times fewer rows. The hour rollup reads 360 times fewer for a week.

    H

    Partition pruning

    from EXPLAIN
    partitions scanned, from EXPLAINWHERE ts >= now - 24 hmetrics_20261002: skipped02metrics_20261003: skipped03metrics_20261004: skipped04metrics_20261005: skipped05metrics_20261006: skipped06metrics_20261007: skipped07metrics_20261008: scanned08metrics_20261009: scanned09metrics_20261010: skipped102 of 9WHERE ts >= now - 7 daysmetrics_20261002: scanned02metrics_20261003: scanned03metrics_20261004: scanned04metrics_20261005: scanned05metrics_20261006: scanned06metrics_20261007: scanned07metrics_20261008: scanned08metrics_20261009: scanned09metrics_20261010: skipped108 of 9WHERE date_bin(1 h, ts) >= now - 24 hmetrics_20261002: scanned02metrics_20261003: scanned03metrics_20261004: scanned04metrics_20261005: scanned05metrics_20261006: scanned06metrics_20261007: scanned07metrics_20261008: scanned08metrics_20261009: scanned09metrics_20261010: scanned109 of 9
    • The planner compares the filter with each partition's bounds and skips the rest.
    • A prepared statement with a generic plan still prunes, when execution starts.
    • Keep the count to a few thousand at most. Planning time and memory grow with it.

    I filter on the partition column itself. A function around it hides the bounds, and the planner scans every partition.

    I

    Approximate answers

    exact or sampled
    7-day averagerows readtimeanswer
    Exact, every row3,024,000491 ms54.9983
    TABLESAMPLE SYSTEM (1)29,6735.8 ms55.0463
    questionapproximate toolstatus
    Total, averageSample of pages or rowsApproved
    Distinct countHyperLogLog sketchApproved
    p99 latencyQuantile sketch or histogramApproved
    Per-bucket counts on a sampleSample of pagesLarge buckets
    Money, billingAny approximationNot approved

    The sample was off by 0.087%. SYSTEM picks whole pages, so rows that arrived together are sampled together.

    For a total or an average I can read a sample. For distinct counts I use a HyperLogLog sketch.

    J

    Column encoding

    delta and run-length, measured
    one day · 50 hosts · 432,000 values per column100 B1 kB10 kB100 kB1 MB10 MBtimestamp, every 10 s: plain 3.46 MB, delta+rle 355 Btimestamp, every 10 sdelta+rle3.46 MB355 Btimestamp, ±1 s jitter: plain 3.46 MB, delta 432 kBtimestamp, ±1 s jitterdelta3.46 MB432 kBhost id: plain 3.46 MB, rle 153 Bhost idrle3.46 MB153 BCPU %, 2 decimals: plain 3.46 MB, delta 810 kBCPU %, 2 decimalsdelta3.46 MB810 kBrequest counter: plain 3.46 MB, delta-of-delta 632 kBrequest counterdelta-of-delta3.46 MB632 kBup flag, 0 or 1: plain 3.46 MB, rle 1.03 kBup flag, 0 or 1rle3.46 MB1.03 kB
    rows
    20 bytes a sample: 8.64 MB for the day
    columns
    time + host + cpu, each with its best scheme: 0.811 MB, 10.7× smaller
    cited
    Prometheus stores 1 to 2 bytes per sample. The Gorilla paper (VLDB 2015) reports 1.37 bytes per point.
    encode one columnpseudo code
    encode(column, scheme):
      v = column
      IF delta:          v = each value minus the one before1
      IF delta-of-delta: v = delta applied twice2
      IF runs: write (value, count) per run of equal values3
      ELSE:    write each value as a varint4    // 1 byte: -64..63
    
      // every 10 s: 1000, 1010, 1020 ...  ->  1000, 10, 10 ...
      //   as runs:  (1000, 1) (10, 8639)5
    1. 1Counters and timestamps turn into small numbers. Small numbers take 1 or 2 bytes as varints.
    2. 2A regular interval becomes zeros. Jitter makes small numbers around zero.
    3. 3Long runs (a host id, an up flag, a fixed interval) collapse to one pair each.
    4. 4Zigzag varint: 7 bits per byte, sign folded into the low bit.
    5. 5A day of timestamps for one host is two pairs. For 50 hosts, 355 bytes in all.
    Tested source Go: encode a column
    Go: encode a columngo
    
    // deltas returns the first value, then each value minus the one before it.
    func deltas(xs []int64) []int64 {
      out := make([]int64, len(xs))
      for i, x := range xs {
        if i == 0 {
          out[i] = x
        } else {
          out[i] = x - xs[i-1]
        }
      }
      return out
    }
    
    // Encode writes a column with one scheme: a count, then the values. A regular timestamp column
    // becomes all zeros after delta-of-delta, and run-length encoding stores those zeros as one pair.
    func Encode(xs []int64, s Scheme) ([]byte, error) {
      out := binary.AppendUvarint(nil, uint64(len(xs)))
      vals := xs
      switch s {
      case Delta, DeltaRLE:
        vals = deltas(xs)
      case DeltaOfDelta, DeltaOfDeltaRL:
        vals = deltas(deltas(xs))
      }
      switch s {
      case Plain:
        for _, v := range vals {
          out = binary.LittleEndian.AppendUint64(out, uint64(v))
        }
      case Varint, Delta, DeltaOfDelta:
        for _, v := range vals {
          out = binary.AppendVarint(out, v) // zigzag: small negatives stay small
        }
      case RLE, DeltaRLE, DeltaOfDeltaRL:
        for i := 0; i < len(vals); {
          run := 1
          for i+run < len(vals) && vals[i+run] == vals[i] {
            run++
          }
          out = binary.AppendVarint(out, vals[i])
          out = binary.AppendUvarint(out, uint64(run))
          i += run
        }
      default:
        return nil, fmt.Errorf("unknown scheme %q", s)
      }
      return out, nil
    }
    
    • Sort the data by the columns you filter on: tenant or host, then time. Runs get long and values close.
    • Pick an encoding per column, then add a general compressor such as LZ4 or ZSTD on top.
    • A jittered timestamp has no long runs: plain delta wins at about 1 byte a value.

    Sorted by host and time, a timestamp column is a few runs, and a gauge needs about 1.9 bytes a value. The three columns take 10.7 times less than rows.

    K

    Failure cases

    what breaks, and how to stay safe
    eventresulthow to stay safesaved by
    A label holds user ids or request ids (high cardinality)One series per value. Memory and the index grow until the metrics server fails.Keep labels to bounded sets: host, route, status. Send per-user data to logs or the warehouse.Bounded labels
    No retention policyRaw data grows every day. Queries and backups slow down, disks fill.Drop raw partitions after N days. Keep rollups longer.DROP partition
    Tomorrow's partition was not createdWrites after midnight fail: no partition for the row.Create partitions days ahead. Alert when fewer than 2 are ahead.Daily job
    A sample arrives after its minute was rolled upThe rollup misses it.Re-read a trailing window each run. The upsert replaces the row.ON CONFLICT
    A filter wraps ts in a functionEvery partition is scanned: 9 of 9.Filter on ts directly. Check the plan.Pruning
    Dashboards scan the OLTP primaryLarge scans compete for CPU and cache with checkout.Read rollups, a replica, or the warehouse.CDC
    The CDC reader stopsThe replication slot keeps WAL on the primary, and its disk fills.Alert on slot lag. Cap kept WAL with max_slot_wal_keep_size.Slot limit
    A chart averages minute averagesWrong when minutes hold different counts.Store sum and count. Divide at the end.Rollup schema
    L

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Postgres table2+ day partitions3+ rollups4+ TimescaleDB5+ CDC, column storemore load →
    Rows one dashboard query reads1k10k100k1M10M7 days, raw rows: 3,240,000 rows7 days, raw rows3,240,00024 h, raw rows: 648,000 rows24 h, raw rows648,0007 days, minute rollup: 540,000 rows7 days, minute rollup540,00024 h, minute rollup: 72,000 rows24 h, minute rollup72,0007 days, hour rollup: 9,000 rows7 days, hour rollup9,000rows, log scale
    stepaddit handlesmove up when you see
    1One table with an index on (host, ts).Recent windows for one host read by index. Millions of rows.Retention deletes rows one by one: 454 ms for one day here, plus dead rows for vacuum.
    2Day partitions and a daily job.Retention by DROP. Windows read only their days: 147 ms for 24 h here.A week-long chart reads every raw row: 922 ms here.
    3Minute and hour rollups, upserted.Week charts in 2.8 ms from the hour rollup. Rollups keep years.Many rollups and compression to maintain by hand.
    4TimescaleDB on the same Postgres.Automatic partitions, continuous aggregates, column compression.Billions of rows, ad hoc scans over many columns, data from many sources.
    5CDC into a column store.Ad hoc SQL over all history, away from the OLTP primary.Top of the ladder. Scale the warehouse itself.

    Steps 1 to 3 are measured in the lab. Do not start at step 5: a warehouse adds a pipeline to run, and dashboards then lag the primary by seconds or minutes.

    I start with one Postgres table and a time index. I add partitions, then rollups. A column store comes when ad hoc scans over all history matter.

    M

    Drill

    predict, then reveal

    0 of 9 known

    1. Why does a column store read less to compute avg(cpu)?

    2. Why does a column of regular timestamps shrink to almost nothing?

    3. Why does the rollup store count and sum, not the average?

    4. A dashboard filters on date_bin('1 hour', ts) and scans all 9 partitions. Why?

    5. How do you delete samples older than 30 days?

    6. When is a materialised view enough, and when do you need a rollup table?

    7. A sample arrives 3 minutes late, after its minute was rolled up. What happens?

    8. A team adds user_id as a Prometheus label. What breaks?

    9. Can an approximate answer be good enough?

    N

    Numbers to say

    measured in the lab
    row
    52.3 bytes per sample in a Postgres heap, for 20 bytes of data.
    columns
    10.7× smaller than rows with delta and run-length encoding.
    rollup
    9× fewer rows per minute bucket; 360× fewer per hour for a week.
    week
    922 ms raw, 2.8 ms from the hour rollup.
    retention
    Drop a day: 0.79 ms. Delete it: 454 ms.
    sample
    1% of pages: 0.087% off, in 5.8 ms.
    Prometheus
    1 to 2 bytes per sample; 15 days kept by default.

    Postgres 16 on an 8-core laptop, one process per query, no parallel workers. Use these as orders of magnitude.