Apache Ozone is a distributed object store from the Hadoop ecosystem. It speaks the Hadoop file system interface, so Spark, Hive and HBase can run on it, and it speaks the S3 protocol, so object-store tools can too. It exists because HDFS has a ceiling: one NameNode holds every file and block in its heap, which caps the number of files a cluster can hold and makes small files expensive.

Ozone's answer is architectural. It splits the namespace from block management, groups blocks into large containers so the metadata that must stay in memory shrinks dramatically, and replicates every metadata service with the Raft protocol. This article walks through each component, follows a write and a read end to end, explains bucket layouts and container states, and ends with operations, failure modes, trade-offs and a checklist. It assumes you know HDFS basics; if not, start with the HDFS NameNode and the HDFS write pipeline. Facts here were checked against the Apache Ozone documentation for release 2.0.0; later 2.x releases exist, so confirm defaults against the version you run.

Advertisement

The problem Ozone solves

In HDFS the NameNode keeps the whole namespace and the block map in memory, and every datanode reports every block it holds. A cluster with hundreds of millions of files needs a very large heap, garbage collection pauses grow, block reports become a significant load, and the file count is effectively capped. Small files make it worse because each one costs a full metadata entry and at least one block entry regardless of size; the small-files problem covers this in detail.

Ozone attacks the problem in two ways. First, it splits the work: the Ozone Manager (OM) owns names, and the Storage Container Manager (SCM) owns where data lives. Second, SCM does not track individual blocks across the cluster. It tracks containers, which are 5 GB by default and hold many blocks. The documentation gives the scale of the effect: a datanode with 196 TB holds around 40,000 containers, compared with about 1.5 million HDFS blocks, a 40-fold reduction in what must be reported.

Components

A cluster has five kinds of service.

  • Ozone Manager holds the namespace of volumes, buckets and keys in RocksDB, including an open-key table for keys being written but not yet committed. It authorises requests and issues block tokens. Three OMs form a Ratis (Raft) group, and only the leader serves writes.
  • Storage Container Manager allocates blocks, creates and tracks containers and pipelines, monitors datanodes, drives re-replication when replicas are lost, and runs the certificate authority for a secure cluster. Its state, including pipelines, containers and deleted blocks, is in RocksDB and is Ratis-replicated across SCM nodes.
  • Datanodes store containers on their disks and serve reads and writes directly to clients. They heartbeat to SCM and send container reports.
  • S3 Gateway translates the S3 REST protocol into Ozone calls. It is stateless, so you run several behind a load balancer.
  • Recon collects data from OM, SCM and datanodes to give a monitoring UI and API, such as missing containers, open keys and capacity use.
Apache Ozone: namespace in Ozone Manager, block space in SCM, data on datanodesClientofs://, o3fs://, ozone shS3 clientaws s3, SDKsS3 Gatewaystateless, scales outOzone Managervolumes, buckets, keys3 OMs, Ratis-replicated RocksDBSCMcontainers, pipelines3 SCMs, Ratis; CA for certificates1 open key2 allocate blockReconmonitoring, insightsDatanode Acontainers on disksDatanode Bcontainers on disksDatanode Ccontainers on disksOne Ratis pipeline: leader + 2 followers replicate each write3 write chunksheartbeats, reports4 commit keyData never flows through OM or SCM; they hand out locations and tokens
The four numbered steps of a key write. OM and SCM hand out names, locations and tokens; bytes travel only between client and datanodes.

Release 2.0.0 removed support for running OM and SCM without Ratis, so high availability is now the only mode of operation; a single-node test cluster is simply a Ratis group of one.

Advertisement

Namespace: volumes, buckets, keys and layouts

The namespace has three levels. A volume is an administrative unit, created by administrators, used for ownership and quota. A bucket lives in a volume and is what users create and point applications at. A key is an object in a bucket. Hadoop clients address paths as ofs://<om-service>/<volume>/<bucket>/<key>, where the ofs scheme exposes all volumes and buckets as one file system tree. S3 clients see buckets in a designated volume.

Each bucket has a layout, which decides how its keys are stored in OM.

LayoutMetadata shapeUse when
FILE_SYSTEM_OPTIMIZED (default)Separate directory and file tables, keyed by parent id plus nameHadoop workloads with directories, renames and deletes of whole trees
OBJECT_STOREFlat keys, full path as the keyPure S3 use, no directory semantics needed
LEGACYOlder flat formatBuckets created before layouts existed

The FSO layout is why Ozone can rename or delete a directory atomically in constant time: the directory is one row, and its children refer to it by id, not by full path. In a flat layout, renaming a directory means rewriting every key beneath it, which is exactly the slow, non-atomic behaviour that makes Spark and Hive job commits painful on plain object stores. The default comes from ozone.default.bucket.layout and is FILE_SYSTEM_OPTIMIZED. Choose per bucket at creation; changing the layout later means copying the data.

Containers, blocks and chunks

Data is stored in a hierarchy. A key consists of one or more blocks, 256 MB by default (ozone.scm.block.size). A block lives inside a container, 5 GB by default (ozone.scm.container.size), and is written as chunks. A block ID is a container ID plus a local ID within the container. To read, a client asks SCM where the container lives, then asks a datanode for the local ID.

Containers are the unit of replication, reporting and recovery. An open container receives writes through a pipeline, a set of datanodes; for the default three-way replication this is a three-node Ratis group with a leader, so every chunk write is ordered and replicated by Raft. When a container fills, or its pipeline fails, SCM moves it from open to closing and then closed. A closed container is immutable, and SCM's replication manager ensures it has the right number of healthy replicas, copying it to a new datanode if one is lost.

If a container must close without its Ratis group agreeing, for example because the pipeline lost a node, a replica can end up quasi-closed, meaning it is closed locally without consensus. SCM then resolves which replicas are authoritative before marking the container closed. Quasi-closed containers in Recon after a node failure are expected; ones that stay that way need investigation.

Following a write and a read

Here is the documented write path, step by step, for a single-block key.

  1. The client asks OM to create a key in a bucket. OM checks permissions and records an entry in the open-key table.
  2. OM asks SCM for a block. SCM picks an open container on a suitable pipeline, three datanodes for replicated data, and returns the block ID.
  3. OM records the block location and returns it to the client with a block token, which proves the client may write that block.
  4. The client streams chunks to the pipeline's datanodes, presenting the token. Ratis replicates each write.
  5. When the data is written, the client calls OM to commit the key with its final length and block list. OM moves it from the open-key table to the key table, and only now does the key become visible.

A read is shorter: the client asks OM for the key's block list and tokens, then reads directly from datanodes, choosing a replica and falling back to another on failure. Neither OM nor SCM is in the data path, which is why their load scales with the number of operations, not with bytes.

The commit step has a practical consequence. Until commit, a key does not exist to readers. A writer that crashes leaves an open key that never commits, and its blocks are cleaned up later, so writers never produce half-visible objects. Ozone 2.0 also added hsync and lease recovery, the semantics HBase needs to persist a write-ahead log before a file is closed.

Worked example: a bucket for Spark and S3

A data platform team wants one store for Spark jobs, which need atomic directory renames for job commits, and for an ingestion service that uses the AWS SDK. They create a volume per business unit and an FSO bucket for the lake.

# Administrator: a volume, then an FSO bucket for the lake
ozone sh volume create /analytics
ozone sh bucket create /analytics/lake --layout FILE_SYSTEM_OPTIMIZED

# Archive bucket with erasure coding instead of 3x replication
ozone sh bucket create /analytics/archive --type EC --replication rs-6-3-1024k

# Hadoop file system view
ozone fs -mkdir -p ofs://ozone1/analytics/lake/events/dt=2026-10-01
ozone fs -put part-0000.parquet ofs://ozone1/analytics/lake/events/dt=2026-10-01/

# Spark reads the same path
spark.read.parquet("ofs://ozone1/analytics/lake/events/")

# S3 clients reach buckets through the S3 gateway endpoint
aws s3 ls --endpoint-url http://s3g.internal:9878 s3://ingest/

Sizing the metadata tier is the main capacity question. Every file and directory is a RocksDB row in OM, replicated to all three OMs, so OM disks should be fast SSDs and sized for the row count. The datanode side is sized like any storage cluster: raw capacity divided by the replication overhead, 3x for Ratis replication or 1.5x for RS-6-3. For the lake bucket, 1 PB of logical data needs about 3 PB raw; the same data in the EC archive needs about 1.5 PB. EC writes cost more CPU and need at least nine datanodes for RS-6-3, one per data and parity block, so keep hot, small and frequently rewritten data replicated. The trade-off is the same as in HDFS erasure coding.

Erasure coding

Ozone supports three built-in EC configurations, written as codec, data blocks, parity blocks and chunk size: RS-3-2-1024k, RS-6-3-1024k and XOR-2-1-1024k, with RS-6-3 recommended. Set EC at bucket creation with --type EC --replication rs-6-3-1024k, or change it for new writes with ozone sh bucket set-replication-config. Data and parity blocks of a block group are spread across distinct datanodes. When a read finds a block missing, the EC client reconstructs the data on the fly from the remaining blocks, a degraded read that is transparent to the application but costs extra network and CPU. SCM's replication manager rebuilds lost blocks in the background.

Operations and failure modes

  • Datanode loss. Open pipelines on the node close, their containers move to closing and SCM creates new pipelines for new writes. Closed containers are re-replicated. Watch under-replicated container counts in Recon until they reach zero.
  • OM leader failover. Clients configured with the OM service ID and all three OM addresses retry against the new leader. Clients configured with one host fail instead. A slow follower forces snapshot installs; keep OM disks fast and similar.
  • Open-key build-up. Jobs that create keys and crash leave open keys and allocated blocks until cleanup runs. Recon reports open keys; a large count usually points at a misbehaving client.
  • Wrong layout. Running Spark commits against an OBJECT_STORE bucket gives you object-store rename behaviour: copies, not atomic moves.
  • Metadata disk exhaustion. OM and SCM RocksDB and Ratis logs grow with operations. Monitor their disks separately from datanode capacity.
  • Security. Secure clusters use Kerberos for users and services, certificates issued by SCM's CA for services, and block tokens for data access. Plan certificate renewal; see Kerberos in Hadoop.
  • Decommissioning. Decommission datanodes through SCM so containers are copied off first; since 2.0 SCM nodes can be decommissioned as well.

Trade-offs against HDFS

AspectHDFSOzone
Metadata scaleBounded by NameNode heap; federation splits itOM in RocksDB; far more objects per cluster
Small filesExpensiveCheaper metadata, still one row per file
ProtocolsHadoop FS, WebHDFSHadoop FS (ofs, o3fs), S3, native client
Directory renameAtomicAtomic in FSO buckets only
Maturity of toolingVery matureYounger; check tool support per version

What to do next

  1. Run a local cluster from the official Docker Compose setup and walk a write through OM, SCM and a datanode in Recon.
  2. Decide layouts up front: FSO for anything Spark, Hive or HBase touches, OBJECT_STORE for pure S3 buckets.
  3. Put OM and SCM metadata on dedicated SSDs and alert on their free space separately.
  4. Configure every client with the OM service ID and all OM addresses, then test a leader failover under load.
  5. Choose replication per bucket: Ratis three-way for hot data, RS-6-3 for cold data when you have enough datanodes.
  6. Track under-replicated containers, quasi-closed containers and open keys in Recon as standing dashboards.
  7. Read the release notes of your 2.x version for changed defaults before upgrading.
Key takeaway: Ozone scales past HDFS by giving names to Ozone Manager, locations to SCM and bytes to datanodes, and by tracking 5 GB containers instead of individual blocks. Every metadata service is Raft-replicated, writes become visible only at commit, and FSO buckets keep the atomic directory semantics Hadoop jobs depend on. Choose layouts and replication per bucket, keep metadata disks fast, and watch container health in Recon.