Performance¶
The performance of a deployment depends on your data, your configuration and your queries. This page tells you what each cost depends on, which configuration keys change it, and how to measure it on your own data. It does not give benchmark results to compare with other systems.
TraceLake has one published target for time: a new event becomes visible to queries in 7 minutes or less for 95 % of events, with the default configuration. All other numbers on this page are examples from specific measurements. They show the size and the direction of an effect. Do not use them as the numbers of your deployment.
Freshness¶
Freshness is the time from the arrival of an event at Kafka to the moment a query can return it. It is the sum of four stages:
| Stage | Configuration key | Default |
|---|---|---|
| The WAL file rotates | ingestion.wal.maxFileAgeMinutes |
5 minutes |
| The file is uploaded | 15 seconds allowed; the uploader looks for files every 5 seconds | |
| The file is committed to the table | ingestion.commitIntervalSeconds |
20 seconds |
| The gateway sees the new snapshot | queryGateway.snapshotPollIntervalSeconds |
8 seconds |
With the defaults, the configured waits of the four stages add up to 343 seconds (300 + 15 for
the upload + 20 + 8). This is not an upper limit. The ingestor examines the age of an idle file
every 30 seconds, each pod changes the commit and poll intervals by a small amount, and the upload
and the commit take time. In the measurement below, the 95th percentile at the defaults was 382
seconds. A WAL file also
rotates when it reaches ingestion.wal.maxFileSizeBytes, so a dataset with a high volume is
fresher than a dataset with a low volume.
The lever is maxFileAgeMinutes, and it costs files. Examples from two measurements:
maxFileAgeMinutes |
Freshness, 95th percentile | Files written in 3 hours by a low-volume dataset |
|---|---|---|
| 5 | 382 s | 37 |
| 2 | 173 s | 89 |
| 1 | 102 s | 175 |
Under continuous high load the key has almost no effect on the number of files, because each file reaches its maximum size before it reaches its maximum age.
More small files mean more work for the compactor and slower queries until the compactor has merged
them. The compactor merges a partition on a backlog tick only when it has more small files than
compaction.backlogFileThreshold. Below that number the files stay until the scheduled pass, and
the scheduled pass does not merge them if their total size is below compaction.minFileSizeMb
(the compactor page gives the full rule).
Measure it with tracelake_freshness_seconds. To find the slow stage, look at
tracelake_ingestor_upload_latency_seconds, tracelake_catalog_commit_latency_seconds and
tracelake_catalog_snapshot_age_seconds.
Ingest¶
What it depends on
- The size of the records and the number of columns. The ingestor decodes each record and maps it to columns.
- The number of open WAL files. There is one open file for each combination of a Kafka
partition and a value of the partition columns. A partition column with many values that occur at
the same time multiplies the open files, up to
ingestion.wal.maxOpenSegmentsPerPartitionfor each Kafka partition. - The WAL volume. The ingestor syncs each file to disk when it rotates, and deletes it after the commit. A slow volume slows both.
- Your storage and your catalog. If uploads or commits are slow, files collect on the WAL
volume, and at
ingestion.wal.diskHighWatermarkPctthe ingestor pauses its Kafka consumers.
How it scales
One dataset on one pod is one pipeline: one consumer, one decode thread, one WAL. Most of its work is serial, so more CPU cores do not make it much faster. To ingest more, add ingestor pods. Kafka gives each pod a share of the partitions, so the number of partitions of the topic is the limit.
Memory
Two things use memory for each open WAL file: the decoded records that wait for the next write
(tracelake_ingestor_wal_buffered_bytes, limited by ingestion.wal.maxBufferedBytes) and the
Parquet row group that is being written (tracelake_ingestor_wal_rowgroup_bytes, not limited).
The first gauge counts less than the real memory. In one measurement it showed 46 % to 57 % of the memory that the process used for that buffer. Thus a limit of 256 MiB is approximately 0.5 to 0.6 GiB of real memory.
When the buffer is at its limit, the ingestor writes smaller batches, and each record costs more CPU. Example: 1.4 to 1.7 microseconds to bind each record in a batch of 1000 rows, and 4.1 to 4.3 microseconds in a batch of one row. The usual cause is a partition column with too many values.
The cost of a smaller maximum file size. Example from one measurement, with
ingestion.wal.maxFileSizeBytes decreased from 64 MiB to 16 MiB:
- approximately 3.4 times the number of data files;
- approximately 2.4 times the number of storage writes for each MiB;
- ingest of 40.7 to 46.8 MB/s from an offered 50 MB/s in that test setup, where the default gave 49.5 MB/s.
Two runs that should have given the same result were 15 % different, so do not read these numbers as exact.
Measure ingest with tracelake_ingestor_fetched_bytes_total, tracelake_ingestor_bytes_total,
tracelake_ingestor_consumer_lag and tracelake_ingestor_wal_active_files.
Queries¶
The time of a query is mostly the time to read data files from your storage. Thus the important question is how many files, and how many row groups in them, the gateway must read.
What removes files before the scan
| A filter on | Is answered by | Notes |
|---|---|---|
| A partition column, or a column with a narrow range of values in each file | The column statistics in the table metadata | No read of a data file or a sidecar. A time range on a date partition is the best case. |
| An identifier, with an equality or a list of values | The bloom filter, if the column is in bloomFilterColumns |
The file must have a bloom-filter sidecar. |
Text, with an equality or LIKE |
The full-text index, if the column is in tantivy.indexingColumns |
Only files that the compactor has merged have this index. |
What the gateway cannot prune
- A range, a
LIKE, a negation or a function on a bloom column. A bloom filter answers only "can this value be here". - A
LIKEpattern with no literal part as long asingestion.tantivy.ngramSize. The gateway then scans the file. - A value longer than
ingestion.tantivy.maxIndexedValueBytesin the column. The gateway scans the file for a filter on that column. - An
ORin which one side cannot be answered. - A file with no sidecar. The gateway scans it. The result is correct, but the query is slower.
Configuration that changes query cost
| Key | Effect |
|---|---|
partitionMapping |
Decides which files a time range or a partition filter can remove. This is the largest effect. A change later applies to new data only and cannot be undone; see the runbook. |
sortColumns |
The compactor sorts each merged file by these columns. An equality lookup on a sorted column then reads a small part of the file. |
compaction.rowGroupSize |
The gateway prunes by row group. Smaller row groups prune more precisely and make the file metadata larger. |
compaction.targetSizeMb |
Larger files mean fewer files to open for a scan, and more temporary disk and memory in the compactor. |
queryGateway.maxWorkerThreads |
The number of threads of a scan. More threads help a scan of many files, and can make other queries slower. Measure it with your queries. |
queryGateway.indexCache.maxBytes |
A sidecar that is not in the cache is read from storage for each query. |
queryGateway.manifestCache.maxBytes |
Table metadata that is not in the cache is read from storage for each query. |
New data is slower to query than old data. Files that the ingestor wrote and the compactor has not merged yet are smaller, are not sorted, and have no full-text index. A query on the last few minutes reads more files for the same number of rows.
Measure queries with these metrics:
tracelake_gateway_query_duration_seconds: the total time.tracelake_gateway_manifest_files_scannedandtracelake_gateway_files_scanned: the number of files after the first pruning pass, and the number that the scan opens. If the two numbers are almost equal, the sidecars do not help the query.tracelake_gateway_prune_candidates_in_totaland_out_total: what each pass keeps.tracelake_gateway_sidecar_miss_total,tracelake_gateway_idx_declines_totalandtracelake_gateway_bloom_unserved_total: why a pass could not prune.tracelake_gateway_datafile_bytes_total: the bytes read from storage.
Gateway memory¶
| What | Limit | Notes |
|---|---|---|
| The query engine | queryGateway.memoryLimit |
No limit if you do not set it. On the HTTP interface the limit applies to each query, so the largest use is memoryLimit × maxInFlightHttpQueries. On the PostgreSQL interface all connections share one limit. |
| Table metadata | queryGateway.manifestCache.maxBytes × (1 + maxInFlightHttpQueries + postgres.maxConnections) |
The cache, plus what each query in progress can hold. 8.5 GiB with the defaults and no PostgreSQL interface. |
| Sidecars | queryGateway.indexCache.maxBytes |
The list of the manifests of a table always stays in memory. Example: approximately 250 bytes for each data file, which is approximately 300 MiB for a table of 1.2 million files. The number of files is thus a cost for the gateway, and compaction decreases it.
Compaction¶
Memory. The compactor keeps the merged files of one partition in memory until it uploads them. The memory that a pass needs grows with the total size of the input files of the largest partition that the pass merges. If a partition is too large, merge it in slices with manual runs.
Temporary disk. Two things use the temporary volume of the pod during a merge:
| Use | Depends on | Example from one measurement |
|---|---|---|
Sorted runs, if the dataset has sortColumns |
The size of the partition after it is written again | 362 MiB for a partition of 3.69 GB in 229 files |
| The full-text index that is being built | One output file and the quantity of text in it | 1.34 GiB for one output file of 128 MB |
The full-text index of a column with long text can be much larger than the data. In one measurement
it was 11 times the size. Put only columns that queries filter on in tantivy.indexingColumns.
tracelake_compactor_ephemeral_peak_bytes shows the largest use of the most recent pass.
One compactor. The compactor runs as one pod. A second pod merges the same partitions again. To merge more at the same time, use manual runs on different datasets, date ranges or partition values.
Measure compaction with tracelake_compactor_duration_seconds,
tracelake_compactor_partition_files_total and tracelake_compactor_sidecar_bytes.
Storage¶
- Compression.
compaction.compressionandcompaction.compressionLevelapply to all data files. A higherzstdlevel makes the files smaller and uses more CPU in the ingestor and the compactor. - The
_unmappedcolumn. Each top-level field of the payload whose name is not the name of a column is stored in_unmapped, also when a mapping reads it. If you map most of the payload to columns with different names, you store most of the payload two times. Give a column the same name as its field to prevent this. - Sidecars. The ratio of the full-text index to the data is
rate(tracelake_compactor_sidecar_bytes_sum{kind="idx"}) / rate(tracelake_compactor_sidecar_bytes_sum{kind="data"}). - Old files. A file that a merge replaced stays in your storage until snapshot expiry and
orphan cleanup remove it. That is
compaction.snapshotRetentionHoursafter the merge, plus the intervals of the two maintenance loops. For that period the merged data is stored two times.
Before you go to production¶
- Send a copy of your real traffic to a test deployment for a minimum of one day, so that the scheduled compaction pass runs at least one time.
- Look at the number of open WAL files and at the number of files in each partition. If either is
high, change
partitionMappingnow. A change later applies to new data only. - Run your real queries, on data of the last minutes and on data that is some days old. Look at the number of files that each query opens.
- Look at the memory of the ingestor, the gateway and the compactor at their highest, and at the temporary disk of the compactor. Set the limits of the pods from those numbers.
- Set the retention of the Kafka topic to the longest outage of your storage or your catalog that you want to survive without data loss.