Azure Cosmos DB is Microsoft's globally distributed, multi-model database service. You write JSON items to containers, it spreads them across partitions and regions, and it promises single-digit-millisecond latency at the 99th percentile within a region. Behind that promise sit a handful of mechanisms that decide whether your application is fast and cheap or throttled and expensive: request units, the partition key, the consistency level, the indexing policy and the change feed.
This article explains each from first principles for the API for NoSQL, the native API, then sizes a real container end to end, shows the Python SDK calls you need, and lists the failure modes that cause most production incidents. Limits and behaviours were checked against Microsoft Learn on 2026-10-03; quotas change, so confirm the numbers that matter to your design before you commit to them.
The architecture in one picture
Request units and throughput modes
Every operation in Cosmos DB is priced in request units, a normalised mix of CPU, IO and memory. The anchor point is that reading one item of about 1 KB by its id and partition key, a point read, costs 1 RU. Larger items, more indexed properties, queries with many predicates or large results, and stronger consistency all cost more. Strong and bounded staleness reads cost about twice as much as weaker levels, because they read from two replicas instead of one.
You buy RUs in one of three ways. Manual provisioned throughput reserves a fixed RU/s per container or database, in steps of 100, with a minimum of 400 RU/s that rises with storage and with the highest throughput ever set. Autoscale sets a maximum Tmax and scales between 0.1 * Tmax and Tmax, billing each hour for the highest level reached, never less than 10 percent. Serverless bills per RU consumed and is limited to a single region. Provisioned RU/s apply in every region of the account, so three regions cost three times the single-region figure.
Every response carries its cost in the x-ms-request-charge header. Treat it as a first-class metric: log it per operation type in development and you will know what your workload costs before it ships.
Logical and physical partitions
A container is split into logical partitions, one per distinct partition key value. All items with the same key live together, which makes the logical partition the unit of transactions: stored procedures, triggers and transactional batches (up to 100 operations) work only within one. A logical partition can hold at most 20 GB.
Logical partitions are hashed onto physical partitions, which you do not control. Each physical partition serves at most 10,000 RU/s and stores at most 50 GB, and is backed by a replica set of four replicas per region. When throughput or data grows past those limits, Cosmos DB splits a physical partition into two; a logical partition is never split. Provisioned throughput is divided evenly across physical partitions, which has a consequence many teams discover late: as storage grows and the partition count rises, each partition's share of the same RU/s shrinks, and a single busy key is capped by its partition's share, never by the container total.
The partition key path is fixed at container creation, and an item's key value cannot be updated in place. Changing either means copying to a new container. Choose a key with high cardinality, even spread of both storage and requests, and alignment with your most common query filter. A query that includes the key in an equality filter goes to one partition; one that does not fans out to every physical partition, and Microsoft documents an overhead of 2 to 3 RU per physical partition checked.
Choosing a key and an indexing policy
When no natural property satisfies all three goals, two techniques help. A synthetic key concatenates properties, for example customerId-yyyyMM, or appends a bucket suffix to spread a hot value, at the cost of queries that must know or enumerate the suffix. A hierarchical partition key declares up to three levels, such as tenant, then user, then session. Queries that filter on the first level are routed only to the partitions holding that prefix, and a first-level value can exceed 20 GB because data is spread by the full path. Global secondary indexes, which copy data into containers keyed differently, are in preview at the time of writing.
from azure.cosmos import CosmosClient, PartitionKey
from azure.identity import DefaultAzureCredential
client = CosmosClient("https://shop-prod.documents.azure.com:443/", DefaultAzureCredential())
db = client.get_database_client("shop")
orders = db.create_container_if_not_exists(
id="orders",
partition_key=PartitionKey(path=["/tenantId", "/customerId"], kind="MultiHash"),
indexing_policy={
"indexingMode": "consistent",
"includedPaths": [{"path": "/status/?"}, {"path": "/createdAt/?"}],
"excludedPaths": [{"path": "/*"}],
"compositeIndexes": [[{"path": "/status", "order": "ascending"},
{"path": "/createdAt", "order": "descending"}]],
},
)The indexing policy in that example is deliberate. By default every property is indexed, which makes ad hoc queries work and makes every write pay for every property. Excluding /* and including only the paths you filter or sort on can cut write cost substantially for wide items. Composite indexes are required for ORDER BY on several properties and speed up filters combined with sorts.
The five consistency levels
Cosmos DB offers five consistency levels, set as an account default and relaxable per request. Within a region, writes always commit to a local majority, three of four replicas; the levels differ in how reads are served and how regions relate.
| Level | Guarantee | Reads from | Writes commit to |
|---|---|---|---|
| Strong | Linearizable: reads see the latest committed write | Two replicas (local minority) | Every region (global majority) |
| Bounded staleness | Lag behind writes bounded by K versions or T time | Two replicas | Local majority |
| Session | Read-your-writes and write-follows-reads within a session | One replica, using the session token | Local majority |
| Consistent prefix | Transactional batches appear together, never out of order | One replica | Local majority |
| Eventual | No ordering guarantee | One replica | Local majority |
Session is the default choice for most applications. After each write the SDK receives a session token for that partition and sends it with later reads, so the replica serving the read must have caught up at least that far. Tokens are partition-bound and cached in the client instance, which matters: a new client, or a different service instance reading data another instance wrote, sees eventual consistency unless you pass the token along. With bounded staleness you configure the allowed lag, and in multi-region accounts it cannot be set tighter than 100,000 writes or 300 seconds; in a single-region account the minimum is 10 writes or 5 seconds. Strong consistency across regions makes each write wait for every region, and it is not available with multiple write regions at all.
Multi-region writes and the change feed
With multiple write regions, two regions can update the same item concurrently. The default last-writer-wins policy keeps the version with the highest _ts, or a numeric property you nominate in the API for NoSQL, and a delete always beats a concurrent insert or replace. The alternative is a custom policy with a merge stored procedure, which can be set only when the container is created; conflicts it cannot resolve go to a conflicts feed your application must drain. Last-writer-wins silently discards one of two concurrent updates, so if two regions can increment the same counter, model the data so that does not happen, for example by giving each region its own item and summing them.
The change feed is an ordered, per-partition log of changes that many designs build on: materialised views, cache invalidation, search indexing and event publishing. Latest version mode, the default, delivers the current version of each created or updated item, can start from the beginning of the container, and does not record deletes, so the usual pattern is a soft-delete flag plus a TTL. All versions and deletes mode records every intermediate change and deletes with metadata, but requires continuous backup, reads only within the backup retention window, and is not supported on accounts that have used partition merge.
Reading, writing and throttling in code
Most of the cost discipline lives in a few calls. Point-read when you know the id and key, scope queries to a partition, and handle throttling explicitly.
from azure.cosmos import exceptions
def get_order(tenant_id, customer_id, order_id):
# Point read: id + full hierarchical key. About 1 RU for a 1 KB item.
return orders.read_item(item=order_id, partition_key=[tenant_id, customer_id])
def open_orders(tenant_id, customer_id):
# In-partition query: routed to one partition, uses the composite index.
return list(orders.query_items(
query="SELECT * FROM o WHERE o.tenantId = @t AND o.customerId = @c AND o.status = @s"
" ORDER BY o.status, o.createdAt DESC",
parameters=[{"name": "@t", "value": tenant_id}, {"name": "@c", "value": customer_id},
{"name": "@s", "value": "open"}],
partition_key=[tenant_id, customer_id],
))
def save(order):
try:
orders.upsert_item(order)
except exceptions.CosmosHttpResponseError as e:
if e.status_code == 429:
# The SDK has already retried; this one escaped. Shed or queue, do not hot-loop.
raise BackPressure(e.headers.get("x-ms-retry-after-ms"))
raise
charge = orders.client_connection.last_response_headers["x-ms-request-charge"]
metrics.observe("cosmos_ru", float(charge), op="upsert")A 429 status means the partition exceeded its share of RU/s in that second. The response carries x-ms-retry-after-ms and the SDKs retry automatically up to a configured limit. A small rate of 429s on a busy container is normal and cheap; sustained 429s concentrated on a few partition key ranges mean a hot partition, which more RU/s fixes only expensively.
Worked example: sizing an order store
An order service for a multi-tenant retail platform expects 1.5 TB of orders in the first year, 4,000 writes per second at peak and 12,000 reads per second. Items average 2 KB. Most reads fetch one order or a customer's open orders.
Key. /tenantId alone has a few hundred values and the largest tenant would pass 20 GB in months. /customerId alone spreads well but loses tenant-scoped queries. A hierarchical key of tenant then customer serves both.
Cost per operation. Measure, do not guess: write 10,000 representative items in a test container with the trimmed indexing policy and read the request charge. Suppose the test shows about 9 RU per write and about 2 RU per point read of a 2 KB item at session consistency. Peak demand is then 4,000 x 9 + 12,000 x 2 = 60,000 RU/s, plus headroom of 30 percent for queries and skew, so autoscale with Tmax of 80,000 RU/s, which floors at 8,000 RU/s overnight.
Partitions. 80,000 RU/s needs at least 8 physical partitions. At 1.5 TB the storage limit of 50 GB needs at least 30. With 30 partitions, each gets about 2,700 RU/s at full scale. The largest tenant's busiest customer must therefore stay well below 2,700 RU/s, which a load test should confirm, and the storage-driven partition count means per-partition throughput will keep falling as data grows unless Tmax rises with it.
Regions. Two regions with a single write region and session consistency give a documented recovery point objective under 15 minutes on regional failure and double the RU bill. If the business needs zero data loss on region failure, that implies strong consistency and its cross-region write latency, a trade-off to settle with product owners explicitly.
Failure modes
- Hot partition. One key or prefix takes most traffic; 429s appear while total consumption looks low. Check normalised RU consumption per partition key range, then fix the key.
- 20 GB wall. A logical partition fills and writes to it fail. Alert well before the limit; the cure is a new key, often hierarchical, and a migration.
- Fan-out queries. Queries without the key hit every physical partition; cost and latency grow with partition count, not with result size.
- Lost session guarantees. A new client instance or a second service reads stale data because it lacks the session token. Use a singleton client and pass tokens across services where read-your-writes matters.
- Index bloat. Default indexing on large items makes writes several times dearer than necessary.
- Throughput dilution. Data growth adds physical partitions, so the same RU/s gives each less. Re-check per-partition headroom quarterly.
Trade-offs
Cosmos DB trades flexibility for predictability. You get documented latency targets, turnkey multi-region replication, five consistency levels and a change feed without operating anything, and in return you design around partition limits, pay per RU in every region, and live with a partition key you cannot change in place. Compared with DynamoDB the shape is familiar, a hashed key-value store with capacity units, but Cosmos DB adds richer queries, tunable consistency and multi-region writes with conflict policies. For relational joins across entities or ad hoc analytics, a relational database or an analytical store fits better.
Related reading: DynamoDB in depth for the closest comparable design, CAP and PACELC for the latency and consistency trade-off behind the five levels, Cassandra partition key design for the same key-choice problem in another system, hot key mitigation for skew patterns, and cloud native databases for the wider landscape.
What to do next
- Write down your top five access patterns and confirm each one either point-reads or filters on the partition key.
- Load-test with realistic key skew and record the request charge per operation type from the response header.
- Trim the indexing policy to the paths you filter and sort on, and add composite indexes for multi-property sorts.
- Compute physical partitions from both RU/s and storage, and check that your busiest key fits in one partition's share.
- Choose the consistency level per request type, use a singleton client, and pass session tokens where read-your-writes crosses services.
- Alert on 429 rate by partition key range, on logical partition size, and on normalised RU consumption.
- Decide between change feed modes before you build on it; all versions and deletes needs continuous backup.