Tutorial · Compose

Add data operators to the chain.

Between the filesystem and your devices, every write passes through a chain of data operators. You choose which operators run, in what order and with what settings, in the same clio.yaml that lists your tiers.

Time
25 minutes
You need
The four-tier setup from Compose local tiers, or any working install
Operators
Erasure coding, replication, placement, tiering, compression

The operator chain

Each operator is a module in its own pool. It receives the storage engine's calls (put, get, size and so on), does its work, and passes them to the pool named by its next_pool_id. The default configuration builds this chain:

data path
filesystem 560.0 → cache 563.0 → indexer 564.0 → [compressor 562.0] → replication 561.0 → core 512.0 → tiers

Block-device operators work one level lower: the storage engine's tiers are block devices, and an operator such as an erasure-coded array is itself a block device built from other block devices.

Two rules for every chain

An entry in compose must come after the entry its next_pool_id names. And when you insert or remove an operator, re-point the next_pool_id of the operator above it, so nothing forwards to a pool that is not there.

OperatorModuleWhat it doesBuilt by default
Erasure-coded arrayclio_safe_bdevStripes data with Reed-Solomon parity over several devices and survives failed devicesYes
Replicationclio_cte_replicationKeeps persistent copies of every write on a durable tierYes
Cacheclio_cte_cacheKeeps node-local raw copies for fast readsYes
Indexerclio_cte_indexerIndexes text for keyword searchYes
Compressorclio_cte_compressorCompresses blobs on the way down, decompresses on the way upNo: -DCLIO_CTE_ENABLE_COMPRESS=ON

Protect data with an erasure-coded array

clio_safe_bdev combines block devices into one array. Data is split into 64 KiB chunks across the data members, and parity members hold Reed-Solomon parity. max_failures is how many members may fail at once without losing data; mark that many members parity: true. With no parity members there is no redundancy.

1. Compose the disks and the array

Add these entries to compose in clio.yaml, above the clio_cte_core entry. Each disk is a file block device on its own drive; here, four drives mounted at /mnt/disk0 to /mnt/disk3.

yaml · clio.yaml
  # Four disks: three hold data, one holds parity.
  - mod_name: clio_bdev
    pool_name: "/mnt/disk0/clio.dat"
    pool_query: local
    pool_id: "330.0"
    bdev_type: file
    capacity: "4GB"
    alloc_log: "/mnt/disk0/clio.alog"
  - mod_name: clio_bdev
    pool_name: "/mnt/disk1/clio.dat"
    pool_query: local
    pool_id: "331.0"
    bdev_type: file
    capacity: "4GB"
    alloc_log: "/mnt/disk1/clio.alog"
  - mod_name: clio_bdev
    pool_name: "/mnt/disk2/clio.dat"
    pool_query: local
    pool_id: "332.0"
    bdev_type: file
    capacity: "4GB"
    alloc_log: "/mnt/disk2/clio.alog"
  - mod_name: clio_bdev
    pool_name: "/mnt/disk3/clio.dat"
    pool_query: local
    pool_id: "333.0"
    bdev_type: file
    capacity: "4GB"
    alloc_log: "/mnt/disk3/clio.alog"

  # One array over the four disks. It survives one failed disk.
  - mod_name: clio_safe_bdev
    pool_name: "raid0"
    pool_query: local
    pool_id: "350.0"
    max_failures: 1
    alloc_log: "/mnt/disk0/raid0.alog"
    members:
      - {pool_name: "/mnt/disk0/clio.dat", pool_id_major: 330}
      - {pool_name: "/mnt/disk1/clio.dat", pool_id_major: 331}
      - {pool_name: "/mnt/disk2/clio.dat", pool_id_major: 332}
      - {pool_name: "/mnt/disk3/clio.dat", pool_id_major: 333, parity: true}

The alloc_log files let the disks and the array recover their layout after a restart.

2. Use the array as a tier

In the clio_cte_core entry, attach the array with existing_pool_id instead of a path and bdev_type. Here it is the persistent tier under a RAM tier:

yaml · clio.yaml
    storage:
      - path: "ram::hot"
        bdev_type: "ram"
        capacity_limit: "4GB"
        score: 1.0
      - path: "raid0"                      # the array as the persistent tier
        existing_pool_id: "350.0"
        existing_pool_module: "clio_safe_bdev"
        score: 0.2                         # matches replication's replica_score
        persistence_level: "long_term"

Give the array a persistence_level other than volatile, or the replication operator will not put its copies there. Its score of 0.2 matches replication's default replica_score. Start the runtime and mount as usual. Files you write are copied to the array, and after clio_run stop and clio_run start they are read back from it.

3. Check its health

shell
$ curl -s "http://127.0.0.1:8080/api/pools/350.0/monitor?query=stats&routing=local"

The reply includes parity_level (1 here), faulty_members, and recovery counters. The dashboard has a page for the array under Pools → clio_safe_bdev, with its members and buttons to add, replace and remove them.

4. Rehearse a disk failure

A file block device has a test hook: while a file named <pool_name>.fail exists, every I/O to that device fails. To see every read go through the array, make the array the only tier for this drill (remove the ram::hot entry and give the array score 1.0).

shell
$ head -c 64M /dev/urandom > data.bin && cp data.bin ~/clio-mnt/
$ touch /mnt/disk1/clio.dat.fail                    # disk1 now fails every I/O
$ cmp data.bin ~/clio-mnt/data.bin && echo intact     # rebuilt from parity
intact
$ curl -s "http://127.0.0.1:8080/api/pools/350.0/monitor?query=stats&routing=local"

Reads still return the right bytes, rebuilt from the other members, and the stats now show "faulty_members": 1. Writes keep working in degraded mode. Replace the failed disk with a new one, and the array rebuilds onto it:

shell
$ sudo mkdir -p /mnt/disk4 && sudo chown $USER /mnt/disk4
$ curl -s -d failed_pool_id=331.0 --data-urlencode member_name=/mnt/disk4/clio.dat \
    -d capacity=4GB -d bdev_type=file \
    http://127.0.0.1:8080/api/mod/clio_safe_bdev/350.0/replace_member
{"ok":true,"failed_pool_id":"331.0","member_name":"/mnt/disk4/clio.dat","member_pool_id":"900.0"}
$ rm /mnt/disk1/clio.dat.fail

Watch recovery_ops_completed catch up with recovery_ops_total in the stats; then faulty_members is back to 0.

Replace a failed disk before you restart

If a member still fails when the runtime starts, the array cannot be created and the runtime does not start. Replace or remove the member while the runtime is running. Afterwards, update the array's entry in clio.yaml to list the new disk.

Keep extra copies with replication

The replication operator is on by default and writes a persistent copy of every write. Its settings, in the clio_cte_replication entry:

KeyDefaultMeaning
replication_factor1Copies across nodes, including the original. Sets remote_copies to one less.
num_replicas0Extra copies on the writing node itself. The default configuration sets 1.
replica_score0.2Where copies go: the tier whose score matches. Below this score, the fast copy may be dropped.
cache_score1.0Where a copy is re-cached when it is read back.
replicate_period_ms50How often pending copies are written. 0 writes them before the write returns.

Over an erasure-coded array, local copies add little: the array already survives a failed disk. In that case you can set num_replicas: 0.

Control how data moves

Three settings in the clio_cte_core entry decide where data goes when it is written and where it moves later.

Placement: where new data lands

yaml
    dpe:
      dpe_type: "max_bw"       # max_bw | round_robin | random

max_bw (the default) picks the fastest tier with room. round_robin rotates across tiers, which suits several devices of the same speed. random is for testing.

Tiering: where data moves later

The data organizer runs periodically and moves blobs whose score no longer matches their tier. Put its keys directly in the clio_cte_core entry, at the same level as storage:

yaml
  - mod_name: clio_cte_core
    pool_name: cte_main
    pool_query: local
    pool_id: "512.0"
    organizer: "frecency"        # none (default) | frecency | hotset | cyclic | scatter | grayscott
    organizer_period_ms: 2000    # how often it runs (default 5000)
    storage:
      # ... your tiers ...

The runtime logs organizer=frecency, period 2000 ms when it starts. frecency scores each blob by how recently and how often it is used, so hot data rises toward RAM and cold data sinks toward disk.

Not under performance

The commented example in clio_default.yaml shows the organizer keys inside the performance: block. They are ignored there. Put them at the top level of the clio_cte_core entry, as above.

Moving data yourself

Any client can give a blob a new score, and the storage engine moves it to the matching tier. This marks a finished run as cold, so its pages move to the slowest tier. Use it with the default organizer: none: an organizer recomputes every score each period and overrides scores you set.

python · cool.py
import os
os.environ["CLIO_WITH_RUNTIME"] = "0"
import clio_cte_core_ext as cte

cte.clio_init(cte.RuntimeMode.kClient, False)
cte.initialize_cte("", cte.PoolQuery.Dynamic())

tag = cte.Tag("/results/run42.bin")           # a file in the mount
for page in tag.GetContainedBlobs():
    if page.isdigit():                         # skip the ~i attribute blob
        tag.ReorganizeBlob(page, 0.1)          # 0.1: the coldest tier
print(tag.GetBlobScore("0"))                   # 0.1

The file stays readable through the mount the whole time. In the dashboard, the cold tier's bytes written grows as the pages move. The cache operator may keep a raw copy in RAM until other data needs the space.

Compress data

Not exercised by the filesystem yet

The compressor picks a codec for each write from fields that the writing client sets. The FUSE mount and the Python API do not set them, so their writes pass through uncompressed. Today, compression applies to C++ clients that request it. The steps below were checked against the source, not run for this page.

The compressor is not in the pip wheel or the release-fuse preset. Build with -DCLIO_CTE_ENABLE_COMPRESS=ON, and use that build for every client too: the option changes a data layout that clients and the runtime share. Codecs are found at configure time, among them zstd, lz4, zlib, lzma, bzip2, snappy and blosc2, plus zfp, sz3 and fpzip through LibPressio.

The default configuration has the compressor entry commented out. Uncomment it, and point the indexer at it so it joins the chain:

yaml
  - mod_name: clio_cte_compressor
    pool_name: clio_cte_compressor
    pool_query: local
    pool_id: "562.0"
    next_pool_id: "561.0"        # replication

  - mod_name: clio_cte_indexer
    # ...
    next_pool_id: "562.0"        # was 561.0

A C++ client requests compression per write through the put's context: dynamic_compress_ = 1 with a codec in compress_lib_ (for example 10 for zstd) and compress_preset_ 1, 2 or 3 (fast, balanced, best), or dynamic_compress_ = 2 to let the compressor choose. A blob is stored compressed only if that saves at least an eighth of its size.

Next: write your own data operator.