A Hadoop cluster is easy to fill and hard to shrink. Storage grows because nothing is ever deleted, compute grows because every new pipeline gets its own queue, and the hardware renewal or the cloud bill arrives as one large number that nobody owns. Cost management is the practice of breaking that number into pieces that teams can see and change, and then pulling the few levers that really reduce it.
This page assumes a running cluster, on premises or in the cloud. Sizing a new cluster is covered in Hadoop capacity planning, and the economics of leaving Hadoop in the cloud migration playbook. Here you will build a rate card, meter compute and storage per tenant with tools the cluster already has, find cold and wasteful data, and apply storage and compute levers, with a worked example and the ways each lever can backfire.
Where the money goes
Start with total cost of ownership for one year, then split it into the two things tenants consume. On premises, that is hardware depreciation, data-centre space and power, network, support contracts or licences, and the operations team. In the cloud, it is instances, block and object storage, data transfer and any managed-service fee. Assign each line to compute or storage. Disks and storage nodes are storage, CPU and memory are compute, and shared items such as the network and staff are split in proportion.
The split matters because the two have different units. Compute is consumed as capacity held over time: a YARN container holding 4 vcores and 16 GB for an hour. Storage is consumed as raw bytes on disk over time, including every replica. Billing tenants for logical bytes hides the most important storage decision, which is the replication or erasure-coding scheme. Billing them for wall-clock job time hides the most important compute decision, which is container size.
A rate card from first principles
A rate card turns annual cost into unit prices: per vcore-hour, per GB-hour of container memory, and per raw TB-month of disk. Divide by usable capacity, not total capacity. A cluster run at 70% average allocation has to recover its whole cost from the 70% it actually allocates, and disks filled to 80% have to recover their cost from that 80%. If you price on full capacity, the showback totals never add up to the bill, and the gap becomes an argument.
HOURS_PER_YEAR = 8760
def rate_card(compute_cost_per_year, storage_cost_per_year,
vcores, memory_gb, raw_tb, target_util=0.70, target_fill=0.80, memory_share=0.5):
"""Unit prices that recover the annual cost at the target utilisation and disk fill."""
per_hour = compute_cost_per_year / HOURS_PER_YEAR
vcore_hour = per_hour * (1 - memory_share) / (vcores * target_util)
gb_hour = per_hour * memory_share / (memory_gb * target_util)
raw_tb_month = storage_cost_per_year / 12 / (raw_tb * target_fill)
return vcore_hour, gb_hour, raw_tb_month
# Example cluster: 120 workers x (40 vcores, 200 GB YARN memory, 12 x 16 TB disks)
print(rate_card(720_000, 480_000, vcores=4_800, memory_gb=24_000, raw_tb=23_040))
# -> about (0.01223, 0.00245, 2.17) per vcore-hour, per GB-hour, per raw TB-monthThe numbers are an example, not a benchmark: a 120-node cluster costing $1.2 million a year, 60% assigned to compute and 40% to storage, with compute split evenly between vcores and memory. Republish the rate card each quarter. Hold it fixed within a quarter, so that tenants see changes in their usage rather than changes in the price.
Metering compute from YARN
The ResourceManager already records what you need. Each application report from the RM REST API includes vcoreSeconds and memorySeconds, the integral of the vcores and megabytes allocated to the application's containers over their lifetime, together with its queue and user. It measures allocation, not utilisation, and that is the right basis for charging. A container that reserves 16 GB and uses 3 GB has still kept 16 GB from everyone else.
import json, urllib.request
from collections import defaultdict
RM = "http://rm.example.internal:8088"
VCORE_HOUR, GB_HOUR = 0.01223, 0.00245 # from the rate card
def finished_apps(begin_ms, end_ms):
url = (f"{RM}/ws/v1/cluster/apps?states=FINISHED,FAILED,KILLED"
f"&finishedTimeBegin={begin_ms}&finishedTimeEnd={end_ms}")
with urllib.request.urlopen(url, timeout=60) as r:
return (json.load(r).get("apps") or {}).get("app", [])
def showback(begin_ms, end_ms):
cost = defaultdict(float)
for a in finished_apps(begin_ms, end_ms):
vcore_h = a["vcoreSeconds"] / 3600
gb_h = a["memorySeconds"] / 1024 / 3600 # memorySeconds is MB-seconds
cost[(a["queue"], a["user"])] += vcore_h * VCORE_HOUR + gb_h * GB_HOUR
return sorted(cost.items(), key=lambda kv: -kv[1])Two operational details matter. The RM keeps only a bounded number of completed applications in memory (yarn.resourcemanager.max-completed-applications), so collect at least hourly and store the results, or read from the timeline or history service. Long-running applications such as Spark streaming jobs and LLAP daemons report their usage so far while still running, so meter them by taking the difference between snapshots. Map queues to cost centres in a small table under version control; see YARN queue hierarchies for designing queues that match your organisation.
Metering storage from HDFS
For each tenant directory you want raw bytes and object count. Raw bytes drive disk cost. Object count, meaning files plus blocks, drives NameNode heap, which is the resource that ends a cluster's growth first when files are small.
# Logical size and raw size (all replicas) per top-level directory
hdfs dfs -du -s -h '/data/*'
# Quotas, remaining quota, directory and file counts per tenant root
hdfs dfs -count -q -h /data/marketing /data/finance /data/archive
# Cap a tenant: 2 million names, 600 TB of raw space (quota counts all replicas)
hdfs dfsadmin -setQuota 2000000 /data/marketing
hdfs dfsadmin -setSpaceQuota 600t /data/marketingIn Hadoop 3, hdfs dfs -du -s prints both the logical size and the space consumed with all replicas. Charge the second. Space quotas count raw bytes too, so a 600 TB quota holds about 200 TB of triple-replicated data. Quotas turn showback into a hard limit. Introduce them only after a few months of reports, with headroom, and with an alert at 80% so jobs do not fail on the first write over the line.
Finding cold and wasteful data
The cheapest terabyte is the one you delete, and the second cheapest is one you move to a cheaper layout. To find candidates, analyse the namespace offline. Fetch the latest fsimage from the NameNode and convert it to a delimited file with the offline image viewer. This puts no load on the live NameNode.
hdfs dfsadmin -fetchImage /var/tmp/fsimage/
hdfs oiv -p Delimited -i /var/tmp/fsimage/fsimage_0000000000123456789 -o /var/tmp/fsimage/ns.tsvThe delimited output has one row per inode, with columns including path, replication, modification and access time, file size and owner. Directories appear too, with a permission string starting with d. A short pandas script turns it into a per-directory report. For namespaces with hundreds of millions of files, load the same file into Hive or Spark instead.
import pandas as pd
df = pd.read_csv("ns.tsv", sep="\t", low_memory=False) # check the first line is the header
files = df[~df["Permission"].str.startswith("d")].copy() # drop directories
files["atime"] = pd.to_datetime(files["AccessTime"], errors="coerce")
files["top"] = files["Path"].str.split("/").str[1:3].str.join("/")
files["raw"] = files["FileSize"] * files["Replication"] # replicated files only; EC differs
files["cold"] = files["atime"] < pd.Timestamp.now() - pd.Timedelta(days=180)
report = files.groupby("top").agg(
files=("Path", "size"),
raw_tb=("raw", lambda s: s.sum() / 1e12),
cold_share=("cold", "mean"),
small_files=("FileSize", lambda s: (s < 16 * 2**20).sum()),
).sort_values("raw_tb", ascending=False)
print(report.head(30).to_string())Access times are only as accurate as dfs.namenode.accesstime.precision, one hour by default. Setting it to 0 disables them, and then every file looks as cold as its modification time. The small-file column matters as much as the cold one. Millions of files under 16 MB cost NameNode heap and task-scheduling overhead far out of proportion to their bytes; the small files problem covers compaction.
Storage levers
- Erasure coding. Three-way replication stores 3 raw bytes per logical byte, while RS-6-3 stores 1.5 and survives the loss of any three blocks in a group. Use it for large, cold files. It writes no extra replicas but costs network and CPU on reads after a failure and on reconstruction, and it removes data locality. Small files are the trap: a file smaller than one cell still gets three parity cells, so it can take more raw space than replication. Policies apply only to new files, so existing data must be rewritten, for example with distcp. See HDFS erasure coding.
- Storage policies and tiering. Tag DataNode volumes as DISK, SSD or ARCHIVE, then set policies such as HOT, WARM or COLD on directories. Setting a policy moves nothing by itself. The
hdfs movertool relocates existing blocks, so schedule it, or the policy only affects new writes. - Compression and file format. Columnar formats such as ORC and Parquet, compressed with ZSTD or Snappy, usually shrink raw text and JSON several times over and speed up scans. Rewriting a log zone once is often the largest single saving.
- Retention. Agree a time to live per dataset class with its owner, and delete by partition. A dataset without a named owner gets a deletion date, not a reprieve.
- Snapshots and trash. Deleting a file that a snapshot still references frees nothing, and deleted files sit in
.Trashforfs.trash.intervalminutes. Audit both before declaring space reclaimed.
# Erasure coding for a cold archive (new files only; rewrite existing ones with distcp)
hdfs ec -enablePolicy -policy RS-6-3-1024k
hdfs ec -setPolicy -path /data/archive -policy RS-6-3-1024k
# Tiering: mark a directory COLD, then move existing blocks to ARCHIVE-tagged disks
hdfs storagepolicies -setStoragePolicy -path /data/logs/2025 -policy COLD
hdfs mover -p /data/logs/2025
# Find space held by snapshots and trash
hdfs lsSnapshottableDir
hdfs dfs -du -s -h '/user/*/.Trash'
Compute levers
Compute waste is mostly the gap between what containers reserve and what they use. Compare allocated memory from the RM with used memory from NodeManager metrics or the job's own counters, and lower mapreduce.map.memory.mb, Spark executor memory and memory overhead where the gap is consistently large. A queue that is half the cluster on paper but busy only at night should have a lower guaranteed capacity and a higher maximum, so it borrows idle capacity instead of reserving it all day; preemption returns the capacity to other queues when they need it. Kill idle applications that still hold containers, such as forgotten notebook sessions. Move batch work with no deadline off-peak, because the peak is what you buy hardware for.
Cloud-specific levers
In the cloud the bill moves with every change, so the levers act faster. Prefer transient clusters that start for a pipeline and terminate after it, keeping data in object storage rather than HDFS. Run task-only nodes, which hold no HDFS blocks, on spot or preemptible capacity, and keep core nodes that store HDFS data on on-demand capacity, because losing them loses blocks. Turn on autoscaling with a sensible minimum, and use object-storage lifecycle rules to move old partitions to colder storage classes. Watch data transfer between regions and to the internet, which YARN metrics never show. Take prices from your provider's current price list.
Worked example: three tenants
Using the example rate card ($0.01223 per vcore-hour, $0.00245 per GB-hour, $2.17 per raw TB-month), the ETL team's queue consumed 1.8 billion vcore-seconds and 9.2 trillion MB-seconds in a month. That is 500,000 vcore-hours and 2.5 million GB-hours, so the compute showback is about $6,100 plus $6,100, roughly $12,200 for the month.
The archive team stores 2.4 PB of logical data at three replicas, 7,200 raw TB, and the showback is about $15,600 a month. The oiv report shows almost all of it unread for over a year, in large files. Rewriting the whole archive under RS-6-3 halves the raw footprint to 3,600 TB and the showback to about $7,800. The marketing team's showback is small, but it owns 40 million files under 1 MB, a fifth of the NameNode's objects. Its action is compaction, not deletion.
Be careful when you report the result. The archive change frees 3.6 PB of raw disk and moves cost between tenants on paper, but the company saves money only when that freed space defers a disk purchase or lets you retire nodes. Report freed capacity and the purchase it avoids next to the showback figures.
Failure modes
- Gaming the meter. Teams that pay for allocation shrink containers until jobs spill to disk or fail. Pair the rate card with an SLA, and review job failure rates along with cost.
- Quotas that break production. A space quota hit at 3 a.m. fails a critical write. Alert well before the quota, and give production pipelines headroom.
- Policies with no movement. COLD set but mover never run, or EC set on a directory whose old files were never rewritten. Verify with the oiv report, not with the configuration.
- Erasure coding on small files. Raw usage rises instead of falling. Filter by file size before converting.
- Phantom reclamation. Deletes held by snapshots or trash. Check
hdfs dfs -duon the parent directory after the trash interval. - Numbers nobody trusts. Showback totals that do not reconcile with the actual bill. Price on target utilisation and publish the method.
What to do next
- Compute last year's total cost and split it into compute and storage. Publish a rate card priced on target utilisation and disk fill.
- Deploy the RM collector hourly, and map every queue to an owner and a cost centre.
- Fetch an fsimage, run the oiv report, and list the top directories by raw bytes, cold share and small-file count.
- Send each owner a monthly showback for three months before introducing any quota.
- Pick one storage lever, usually erasure coding for one large cold archive, and measure the freed raw TB with oiv afterwards.
- Compare allocated with used container memory for the ten most expensive jobs, and right-size them.
- Record freed capacity against the next hardware purchase or cloud budget, so the savings are real.