TT Lab
Get started
Learn Learning paths Courses

ClickHouse — A Columnar Analytics Database from the Inside

Filters outside the sort key — skipping with indexes and projections

Continue in TT Lab

Goal

On a log table whose sorting key is ts, you reduce queries that filter by columns outside the key using skipping indexes and projections, and confirm with EXPLAIN, rows_read and the system tables when each works and when it does not, and what it costs.

Why it matters

Unlike its name suggests, a skipping index is not a B-tree but a summary of bundles of granules. If the data distribution does not fit, it cannot skip a single granule, and on existing parts it has no effect until you MATERIALIZE. A projection lets you use another sort order without changing queries, but it uses a lot of disk. This lab's grader does not trust the numbers you wrote — it re-runs the SELECTs you saved read-only, with the condition cache off, and also runs them with the index turned off, to confirm that the reduced rows are really thanks to that index or projection.

Steps

  1. Create the database skp and the table skp.logs — columns ts DateTime, seq UInt64, service LowCardinality(String), user_id UInt32, trace_id UInt64, latency_ms UInt32, path String (in this order), MergeTree, ORDER BY ts. Insert 1 million rows once with /opt/lab/fixtures/skip/logs.sql and make it one part with OPTIMIZE TABLE skp.logs FINAL.
  2. Run ALTER TABLE skp.logs ADD INDEX idx_trace trace_id TYPE bloom_filter GRANULARITY 1 (do not MATERIALIZE yet), and save the output of EXPLAIN indexes = 1 SELECT count(), sum(latency_ms) FROM skp.logs WHERE trace_id = 12249378055674284861 to /root/ch/skip/explain_before.txt.
  3. Run MATERIALIZE INDEX idx_trace with mutations_sync = 2, save the SELECT above as /root/ch/skip/q_trace.sql, and write the rows_read measured with the condition cache off and that mutation's mutation_id to /root/ch/skip/trace.json.
  4. Add minmax indexes (GRANULARITY 1) named idx_seq on seq and idx_user on user_id, and MATERIALIZE them. Create /root/ch/skip/q_seq.sql, which gives the sum(latency_ms) of rows with seq BETWEEN 1001500000 AND 1001520000, and /root/ch/skip/q_user_range.sql, which gives the sum(latency_ms) of rows with user_id BETWEEN 100 AND 120, and write the rows_read of the two queries to /root/ch/skip/minmax.json as seq_rows_read and user_rows_read.
  5. Add the projection p_user (SELECT * ORDER BY user_id) and MATERIALIZE it. Create /root/ch/skip/q_user.sql, which gives the count(), sum(latency_ms) of rows with user_id = 4242, and write the rows_read with the projection used and the rows_read with it turned off by --optimize_use_projections 0 to /root/ch/skip/proj.json as with_projection and without_projection.
  6. Run the query of step 5 with SETTINGS log_comment = 'chs-skip-06' attached, then from the end (QueryFinish) record in system.query_log write query_id, read_rows and projections to /root/ch/skip/qlog.json.
  7. Add the aggregate projection p_svc_hour (SELECT service, toStartOfHour(ts), count(), avg(latency_ms) GROUP BY service, toStartOfHour(ts)) and MATERIALIZE it. Create /root/ch/skip/q_svc.sql, which gives service, count(), round(avg(latency_ms), 3) per service in service order, and write its rows_read and the rows of p_svc_hour in system.projection_parts to /root/ch/skip/agg.json as rows_read and projection_rows.
  8. Write the bytes_on_disk of the active part (part_bytes), the bytes_on_disk of the two projection parts (p_user_bytes and p_svc_hour_bytes), and the data_compressed_bytes of the two indexes (idx_trace_bytes and idx_seq_bytes) to /root/ch/skip/cost.json.

Notes

Create a log table sorted by time

Create the database skp and the table skp.logs. The columns are, in order, ts DateTime, seq UInt64, service LowCardinality(String), user_id UInt32, trace_id UInt64, latency_ms UInt32, path String, it is MergeTree, and ORDER BY ts. Insert 1 million rows once with /opt/lab/fixtures/skip/logs.sql and make it one part with OPTIMIZE TABLE skp.logs FINAL.

If you merge the parts into one, the number of granules is fixed at one set and you can compare before and after the index on the same scale. One million rows is 123 granules of 8192 rows. seq grows with time, user_id is evenly spread regardless of time, and trace_id differs for every row.

EXPLAIN when you have only defined the index

Run ALTER TABLE skp.logs ADD INDEX idx_trace trace_id TYPE bloom_filter GRANULARITY 1 (do not yet MATERIALIZE). Then save the output of EXPLAIN indexes = 1 SELECT count(), sum(latency_ms) FROM skp.logs WHERE trace_id = 12249378055674284861 to /root/ch/skip/explain_before.txt.

Under Indexes in the EXPLAIN output there is a PrimaryKey slot and a Skip slot. The Granules of the Skip slot is remaining/total. Existing parts do not have an index file yet, so the index filters nothing — also check the size in system.data_skipping_indices.

Rows read after MATERIALIZE

Run ALTER TABLE skp.logs MATERIALIZE INDEX idx_trace SETTINGS mutations_sync = 2, and save the SELECT of step 2 (without EXPLAIN) as /root/ch/skip/q_trace.sql. Write the statistics.rows_read measured with the condition cache off as rows_read, and the id of that mutation found in system.mutations as mutation_id, to /root/ch/skip/trace.json.

MATERIALIZE INDEX is a mutation that writes the index file anew in all parts — notice the mutation number attached at the end of the part name. A Bloom filter discards only granules it can say for certain are absent, so even if the row you are looking for is one line, it may read a few more false-positive granules.

The same minmax, different effects

Add minmax indexes (GRANULARITY 1) named idx_seq on seq and idx_user on user_id, and MATERIALIZE both. Create /root/ch/skip/q_seq.sql, which gives the sum(latency_ms) of rows with seq BETWEEN 1001500000 AND 1001520000, and /root/ch/skip/q_user_range.sql, which gives the sum(latency_ms) of rows with user_id BETWEEN 100 AND 120, and write the rows_read of the two queries to /root/ch/skip/minmax.json as seq_rows_read and user_rows_read.

minmax records only the minimum and maximum per granule. For a column that grows together with the sorting key (ts), each granule's range is narrow and most granules outside the condition drop out, while for a column evenly spread regardless of time, every granule's range is nearly the whole, so not one drops out.

A hidden copy in another sort order

Run ALTER TABLE skp.logs ADD PROJECTION p_user (SELECT * ORDER BY user_id) and then MATERIALIZE PROJECTION p_user with mutations_sync = 2. Create /root/ch/skip/q_user.sql, which gives the count(), sum(latency_ms) of rows with user_id = 4242, and write the rows_read measured with the condition cache off as with_projection, and the value measured with --optimize_use_projections 0 added as without_projection, to /root/ch/skip/proj.json.

A projection is a copy that goes into a subdirectory per part. Leave the query with the original table name as it is, and the optimizer picks the side with the least to read. In the copy sorted by user_id, user_id is the sorting key, so the primary index narrows it to a single granule.

See in the record which projection was chosen

Run SELECT count(), sum(latency_ms) FROM skp.logs WHERE user_id = 4242 SETTINGS log_comment = 'chs-skip-06', and from the end (type = 'QueryFinish') record in system.query_log write query_id, read_rows and projections to /root/ch/skip/qlog.json.

The query statement has no projection name, so what was used can be known only from the record. The projections column of query_log leaves the projections used as an array in the form database.table.name. The log is flushed about once a second, so first SYSTEM FLUSH LOGS.

A pre-aggregated projection

Add the projection p_svc_hour (SELECT service, toStartOfHour(ts), count(), avg(latency_ms) GROUP BY service, toStartOfHour(ts)) and MATERIALIZE it. Create /root/ch/skip/q_svc.sql, which gives service, count(), round(avg(latency_ms), 3) per service in service order, and write the rows_read with the condition cache off and the rows of the active part of p_svc_hour in system.projection_parts to /root/ch/skip/agg.json as rows_read and projection_rows.

A projection with a GROUP BY becomes a hidden AggregatingMergeTree that holds one line of intermediate state for count and avg per hour and service. The per-service aggregation only has to merge those lines again, so it reads a few thousand rows of the copy instead of the 1 million original rows.

The disk price of the things that made it fast

Write the bytes_on_disk of the active part of skp.logs as part_bytes, the bytes_on_disk of the active parts of p_user and p_svc_hour in system.projection_parts as p_user_bytes and p_svc_hour_bytes, and the data_compressed_bytes of idx_trace and idx_seq in system.data_skipping_indices as idx_trace_bytes and idx_seq_bytes, to /root/ch/skip/cost.json.

A part's bytes_on_disk includes the projection subdirectories and index files inside it. If you put side by side a copy with everything re-sorted, an aggregate copy of a few thousand rows, a Bloom filter and a minmax, you can see at a glance which is expensive.