TT Lab
Get started
Learn Learning paths Courses

ClickHouse — A Columnar Analytics Database from the Inside

Load Two Million Rows and Look Inside the Column Files

Continue in TT Lab

Goal

You insert 2 million rows into a MergeTree table and read parts, granules and per-column compression ratios from the system tables. You confirm with rows_read and bytes_read that even when you read the same number of rows, the cost differs depending on the columns touched and the sorting key.

Why it matters

In a column-oriented database, query cost is decided not by "how many rows did it pick" but by "how many granules of which columns did it read". Without this sense, you choose the sorting key wrongly, use SELECT * out of habit, and plan disk capacity by the original size. This lab's grader does not trust the numbers you wrote — it reads the server's system.parts, system.columns and system.query_log directly, re-runs the SELECTs you saved read-only, and measures the rows and bytes read to compare. The grader's own queries are not left in query_log.

Steps

  1. Create the database col and the table col.events — columns ts DateTime, site LowCardinality(String), user_id UInt64, url String, dur_ms UInt32, country FixedString(2) (in this order), engine MergeTree, sorting key ORDER BY (site, ts).
  2. Run /opt/lab/fixtures/columnar/events.sql once to insert 2 million rows.
  3. After merging the parts into one with OPTIMIZE TABLE col.events FINAL, write that part's part (name), rows, marks and part_type to /root/ch/columnar/parts.json.
  4. From system.columns, write the column with the largest uncompressed/compressed ratio (most_compressed), the column with the smallest (least_compressed) and the sum of compressed bytes (total_compressed_bytes) to /root/ch/columnar/columns.json.
  5. Create /root/ch/columnar/q_one.sql, which sums only dur_ms, and /root/ch/columnar/q_url.sql, which reads the contents of the url string (for example max(url)), and write the bytes_read of the two queries to /root/ch/columnar/bytes.json as one_col_bytes and url_bytes.
  6. Create /root/ch/columnar/q_key.sql, which sums dur_ms for rows with site = 'docs.example', and /root/ch/columnar/q_nokey.sql, which sums dur_ms for rows with country = 'KR', and write the rows_read of the two queries to /root/ch/columnar/rows.json as key_rows_read and nokey_rows_read.
  7. Run SELECT uniqExact(user_id) FROM col.events WHERE site = 'shop.example' with SETTINGS log_comment = 'chs-col-07' attached, then find that query in system.query_log and write query_id, read_rows, read_bytes and result_rows to /root/ch/columnar/qlog.json.
  8. Create a table col.small with the same columns as col.events, insert 1000 rows from col.events with a single INSERT, and write the part format of the two tables to /root/ch/columnar/part_types.json as events and small.

Notes

Create a MergeTree table with a sorting key

Create the database col and the table col.events. The columns are, in order, ts DateTime, site LowCardinality(String), user_id UInt64, url String, dur_ms UInt32, country FixedString(2), the engine is MergeTree, and the sorting key is ORDER BY (site, ts).

It is two statements, CREATE DATABASE and CREATE TABLE. A MergeTree without a sorting key cannot be created — because the key is the row order on disk. If you write a column name or type differently by even one character, the source script in a later step will not load.

Insert the 2 million source rows

Run /opt/lab/fixtures/columnar/events.sql once to insert 2,000,000 rows into col.events.

Just pass it to clickhouse-client with --queries-file. This script makes row numbers with numbers() and draws values with a hash, so the same rows come out however many times you run it. If the row count is 4 million, you inserted twice.

Merge the parts into one and read the name

After merging the parts into one with OPTIMIZE TABLE col.events FINAL, write the part (name), rows, marks and part_type of that active part from system.parts to /root/ch/columnar/parts.json.

A single INSERT can create several parts depending on block size. After merging, the name changes (the merge level goes up). Look only at rows with active = 1 — the old parts that were merged stay in the list for a while. If you alias the name column of system.parts as part and receive it as JSONEachRow, you can save it as it is.

How different is the compression ratio per column

From system.columns, compute data_uncompressed_bytes / data_compressed_bytes for each column of col.events, and write the largest column as most_compressed, the smallest as least_compressed and the sum of compressed bytes as total_compressed_bytes to /root/ch/columnar/columns.json.

Using argMax(name, ratio) and argMin(name, ratio), you finish in one query. The column that shrinks best is usually one at the front of the sorting key or one with only a few distinct values, and the one that shrinks least is a value that keeps growing or one close to random.

Same rows, different columns — bytes read

Create /root/ch/columnar/q_one.sql, which sums only dur_ms, and /root/ch/columnar/q_url.sql, which reads the contents of the url string, run both queries with --format JSON, and write statistics.bytes_read to /root/ch/columnar/bytes.json as one_col_bytes and url_bytes.

Both queries read all 2 million rows. What differs is the size of the columns touched. One UInt32 column is 4 bytes per row. On the url side you must not use length(url) — that function is rewritten to read a subcolumn holding only the length and does not read the contents. Use a function that compares contents, like max(url).

Skipping happens only when filtering by the sorting key

Create /root/ch/columnar/q_key.sql, which sums dur_ms for rows with site = 'docs.example', and /root/ch/columnar/q_nokey.sql, which sums dur_ms for rows with country = 'KR', and write the statistics.rows_read of the two queries to /root/ch/columnar/rows.json as key_rows_read and nokey_rows_read.

site is the first column of the sorting key, so the index skips granules that cannot hold that value. country is not in the key, so there is no basis to skip. Think about why the number of rows read comes out near a multiple of 8192 — the unit of skipping is not the row but the granule.

Find your own query in query_log

Run SELECT uniqExact(user_id) FROM col.events WHERE site = 'shop.example' SETTINGS log_comment = 'chs-col-07', find the record of the end of that query (type = 'QueryFinish') in system.query_log, and write query_id, read_rows, read_bytes and result_rows to /root/ch/columnar/qlog.json.

query_log is flushed to disk about once a second. To find it right away, first run SYSTEM FLUSH LOGS. A query leaves two lines, start (QueryStart) and end (QueryFinish), and the number of rows read is only on the end line. If you filter by log_comment, only yours remains among the many queries.

A small part becomes Compact

Create a table col.small with the same columns and engine as col.events, insert 1000 rows from col.events with a single INSERT, and write the active part format (part_type) of the two tables to /root/ch/columnar/part_types.json as events and small.

CREATE TABLE ... AS another_table copies the columns and engine as they are. The part format is Compact if the part's bytes and row count are smaller than the table settings min_bytes_for_wide_part and min_rows_for_wide_part, and Wide if larger. Do not change the table settings.