Skip to main content
Version: 3.0 (next)

Azure Cosmos DB Azure Cosmos DB Integration Guide

Connect to Azure Cosmos DB to read and write items from your pipelines, and to run a pipeline whenever an item is created or updated. This guide covers connection setup, the eight function types, and pipeline integration.

Overview​

Azure Cosmos DB is Microsoft's globally distributed NoSQL database. This connector speaks its native NoSQL (Core SQL) API. To reach a Cosmos DB account through its MongoDB API instead, use the MongoDB connector.

The connector provides:

  • Point reads of one item by its ID and partition key, the cheapest operation Cosmos DB offers
  • Bulk reads of up to 1,000 items in one call, with the missing ones listed
  • SQL queries with named parameters, a limit sent to the service, and continuation paging
  • Writes — create, upsert and replace — with optimistic concurrency by ETag
  • Partial updates with patch, including server-side increments
  • Transactional batches of up to 100 operations on one partition, all or nothing
  • A change feed trigger that runs a pipeline for each created or updated item
  • The Request Units each call consumed, on every result
  • Account key or Microsoft Entra ID authentication
  • The Cosmos DB emulator for local development

One connection, one database​

A connection addresses one database in one account. Every function names the container it works with, so a single connection covers every container in the database. To work with a second database, create a second connection.

Partition keys​

Every Cosmos DB container is split by a partition key — a path such as /deviceId chosen when the container was created. An item is identified by its id and its partition key value together: two items in different partitions can share an id.

Functions that address one item take a Partition Key field. Its value follows one rule:

You writePartition key value
press-1the string "press-1"
42the string "42"
[42]the number 42
[true]the boolean true
[null]null
["ankara", "line-3"]a hierarchical key with two levels

A plain value is always a string key, because a template renders to text and cannot say otherwise. Anything else is written as a JSON array. Every result renders partition keys the same way, so a value read from one node's result passes straight into the next node's Partition Key field.

A string and a number are different partitions

An item stored with "n": 42 lives in partition [42], not 42. Azure Cosmos DB finds nothing when you read it with the string key. The emulator is more lenient and finds it either way — so test numeric keys against Azure before relying on them.

Connection Configuration​

Creating an Azure Cosmos DB Connection​

Navigate to Connections → New Connection → Azure Cosmos DB and configure the following.

1. Profile Information​

FieldDefaultDescription
Profile Name-A descriptive name for this connection profile (required, max 100 characters)
Description-Optional description for this connection

2. Account and Database (Connection tab)​

FieldDefaultDescription
Account Endpoint-The account URI from the Azure portal's Keys or Overview blade, e.g. https://my-account.documents.azure.com:443/ – required
Database-The database this connection addresses – required

3. Authentication (Security tab)​

FieldDefaultDescription
Authentication Methodaccount_keyaccount_key, service_principal, managed_identity or default_credential

Account key (Only displayed when the method is account_key)

FieldDefaultDescription
Account Key-The primary or secondary key from the account's Keys blade. Stored encrypted

A read-only key allows reads, queries and the change feed. Writes then fail with 401 Unauthorized.

Microsoft Entra ID (Only displayed for the three Entra ID methods)

FieldDefaultDescription
Tenant ID-The directory ID. Required for service_principal
Client ID-The app registration's client ID. For managed_identity, the optional client ID of a user-assigned identity
Client Secret-The app registration's secret. service_principal only. Stored encrypted

managed_identity uses the identity of the Azure host MaestroHub runs on. default_credential tries the environment, workload identity, managed identity and the Azure CLI in turn.

Entra ID needs a data-plane role

Azure roles such as Contributor manage the account but do not grant access to its data. Assign a Cosmos DB data-plane role to the identity:

  • Cosmos DB Built-in Data Reader — reads, queries and the change feed
  • Cosmos DB Built-in Data Contributor — everything this connector does

Data-plane roles are assigned with the Azure CLI (az cosmosdb sql role assignment create), not in the portal's Access control blade.

4. Advanced (Advanced tab)​

FieldDefaultDescription
Preferred Regions-Comma-separated Azure regions to read from, nearest first — West Europe, North Europe. Only matters for an account replicated to several regions
Connect Timeout30sBound on connecting and on each health check (1s–300s). Each operation has its own Timeout

How the health check works​

Test Connection — and every later health check — reads the database. That one call proves the endpoint answers, the credentials are accepted, and the database exists. It needs only metadata read permission, which every Cosmos DB data-plane role includes.

A database that does not exist fails the check with its name in the message.

Local development with the emulator​

The Cosmos DB emulator runs the NoSQL API in a container, on Linux, macOS and Windows:

docker run -d -p 8081:8081 \
mcr.microsoft.com/cosmosdb/linux/azure-cosmos-emulator:vnext-preview --protocol http

Then connect with:

FieldValue
Account Endpointhttp://localhost:8081/
Authentication Methodaccount_key
Account KeyC2y6yDjf5/R+ob0N8A7Cgv30VRDJIWEHLM+4QDU5DE2nQ9nDuVTqobD4b8mGGyPMbIZnqyMsEcaGQy67XIw/Jw==

That key is the emulator's published, fixed key — the same on every install, not a secret.

The emulator is not Azure

The emulator accepts any account key and does not check that a partition key header matches the item. The connector checks the item's partition key itself, so writes behave the same on both — but a wrong key only fails against Azure.

Function Builder​

Creating Azure Cosmos DB Functions​

Open the connection, go to the Functions tab, and click New Function. Pick one of the eight operations.

Selecting an Azure Cosmos DB function type

The eight Azure Cosmos DB function types

Items in and out​

Items are plain JSON in both directions. On the way out, the connector:

  • delivers a whole number as an integer, keeping every digit — a 19-digit identifier is not rounded through a float
  • removes the resource-link properties Cosmos DB adds to every item (_rid, _self, _attachments, _lsn), which address the item inside the service's own REST model and mean nothing to a pipeline
  • keeps _etag, which the If-Match field reads, and _ts, the time the item last changed (epoch seconds)

Every result also carries requestCharge: the Request Units (RU) the call consumed. RU are what Cosmos DB bills and what its throughput limit counts. A container over its provisioned RU/s answers 429; the connector retries it up to three times, waiting as long as the service asks, and after that reports it as a retryable failure.

Read Item​

Read one item by its ID and partition key.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container holding the item
Item IDstringyes-The item's id
Partition Keystringyes-The item's partition key — see Partition keys
Timeoutdurationno30mBound on this operation

A point read is the cheapest operation Cosmos DB has: one RU for a 1 KB item. An item that does not exist is not an error — the function succeeds with found false and a null item. A container that does not exist is an error, naming the container.

Container:     machines
Item ID: ((deviceId))
Partition Key: ((deviceId))

Read Many Items​

Read up to 1,000 items by ID and partition key in one call.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container holding the items
ItemsJSONyes-Array of {"id": …, "partitionKey": …}, at most 1,000
Timeoutdurationno30mBound on this operation
[
{"id": "press-1", "partitionKey": "press-1"},
{"id": "press-2", "partitionKey": "press-2"},
{"id": "line-3", "partitionKey": ["ankara", "line-3"]}
]

Cosmos DB groups the items by partition and fetches each group in one request, which costs far less than a Read Item per entry in a loop. Inside the array, partitionKey can be a JSON string, number, boolean, null, or an array for a hierarchical key.

Items that do not exist come back in missing, each as {id, partitionKey}. An ID under the wrong partition key counts as missing — it names a different item.

Query Items​

Run a Cosmos DB SQL query.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container to query
Querystringyes-A Cosmos DB SQL query, e.g. SELECT * FROM c WHERE c.deviceId = @device
Query ParametersJSONno-Object of parameters — {"@device": "press-1", "@min": 80}
Partition Keystringno-Run the query inside this one partition. Leave empty to query every partition
Limitintegerno100Maximum items to return, 1–10,000
Continuation Tokenstringno-Resume where a previous run stopped
Timeoutdurationno30mBound on this operation

Reference values as @name parameters rather than splicing them into the query text. A parameter keeps a value from being read as SQL, and the service can reuse the query plan.

Query:            SELECT * FROM c WHERE c.deviceId = @device AND c.celsius > @limit
Query Parameters: {"@device": "((deviceId))", "@limit": 80}

The Limit is sent to Cosmos DB as the page size, so it bounds what is read and billed. A limit above 1,000 is filled over several requests, each asking only for what is still missing.

Paging​

When more results may follow, the result carries a continuationToken and truncated is true. Feed the token back as the next run's Continuation Token, with the same query and partition key:

run 1: limit 100                       → continuationToken "…", truncated true
run 2: continuationToken from run 1 → continuationToken "…", truncated true
run 3: continuationToken from run 2 → truncated false, no token

Cross-partition queries​

Leaving Partition Key empty runs the query across every partition. Scope it to one partition whenever you can: it is cheaper, and some SQL features need it.

ORDER BY, TOP, DISTINCT, GROUP BY and aggregates need a partition key

The Cosmos DB Go SDK this connector uses sends queries through the Azure gateway without a client-side query engine. Across partitions, the gateway runs only simple filters and projections; a query using ORDER BY, TOP, OFFSET … LIMIT, DISTINCT, GROUP BY or an aggregate such as COUNT fails with an error that tells you to set Partition Key. Inside one partition they all work.

Write Item​

Create, upsert or replace a whole item.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container to write into
ItemJSONyes-The item. It must carry a string id and the container's partition key field
Modeenumnoupsertupsert, create or replace
If-Match ETagstringno-Write only if the stored item's _etag still equals this value
Timeoutdurationno30mBound on this operation
ModeWhen the ID existsWhen it does not
upsertreplaces the itemcreates it
createfails with 409 Conflictcreates it
replacereplaces the itemfails with 404 Not Found

The partition key is read from the item's own field — for a container partitioned on /deviceId, from item.deviceId. An item without that field is refused before anything is sent, naming the path. The result says whether the write created a new item (created true) or replaced one.

{
"id": "((readingId))",
"deviceId": "((deviceId))",
"celsius": ((value)),
"recordedAt": "((timestamp))"
}

The item is sent as you wrote it, so a large integer is not rounded on the way in either.

Optimistic concurrency​

Read an item, keep its _etag, and pass it as If-Match ETag on the write. If anyone changed the item in between, the write fails with 412 Precondition Failed and nothing changes. Create ignores If-Match: there is no stored item to compare against.

Patch Item​

Change parts of an existing item without sending the whole of it.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container holding the item
Item IDstringyes-The id of the item to patch
Partition Keystringyes-The item's partition key
OperationsJSONyes-Array of up to 10 patch operations
Conditionstringno-Apply the patch only when the stored item matches, e.g. FROM c WHERE c.status = 'idle'
If-Match ETagstringno-Patch only if the stored item's _etag still equals this value
Timeoutdurationno30mBound on this operation

Each operation is {"op", "path", "value"}, and the path is a JSON Pointer:

opWhat it does
setSets a field, creating it if needed
addAdds a field, or inserts into an array — /alarms/- appends
replaceReplaces a field that must already exist
removeRemoves a field. Takes no value
incrAdds a whole number to a numeric field, on the server
[
{"op": "incr", "path": "/cycles", "value": 1},
{"op": "set", "path": "/status", "value": "((status))"},
{"op": "add", "path": "/alarms/-", "value": "((alarmCode))"}
]

The operations apply atomically, and the result carries the item as it is after the patch — the way to read a counter's new value after an increment.

Two limits come from the Cosmos DB Go SDK, and the form refuses both with the fix spelled out:

  • A value cannot be null. The SDK drops a null value from the request. Use remove to clear a field.
  • Write strings in a Condition with single quotes. The SDK does not escape the condition, so a double quote or backslash would break the request.

Delete Item​

Delete one item.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container holding the item
Item IDstringyes-The id of the item to delete
Partition Keystringyes-The item's partition key
If-Match ETagstringno-Delete only if the stored item's _etag still equals this value
Timeoutdurationno30mBound on this operation

Deleting an item that is not there succeeds with deleted false, so the operation is safe to replay.

Transactional Batch​

Run up to 100 operations on one partition as a single transaction.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container the batch writes into
Partition Keystringyes-The partition every operation belongs to
OperationsJSONyes-Array of up to 100 operations
Timeoutdurationno30mBound on this operation

Each operation has an op and what that op needs:

opNeeds
create, upsertitem
replaceitem (and optionally id, which must match the item's)
patchid and patch — the same operations Patch Item takes
read, deleteid

Any operation can carry ifMatch.

[
{"op": "create", "item": {"id": "((readingId))", "deviceId": "press-1", "celsius": ((value))}},
{"op": "patch", "id": "press-1", "patch": [{"op": "incr", "path": "/readings", "value": 1}]},
{"op": "read", "id": "press-1"}
]

Either every operation applies or none does. When one fails, the whole batch is rolled back and the function fails with the operation that caused it — operations[1] (create "r-1") failed with status 409 Conflict. Every item a batch writes must belong to the batch's partition; one that does not is refused before anything is sent.

Change Feed​

Run a pipeline for each item created or updated in a container. This function is used on an Azure Cosmos DB trigger node, not as a pipeline step.

FieldTypeRequiredDefaultDescription
Containerstringyes-Name of the container to follow
Start Fromenumnonownow fires only for changes after the trigger starts; beginning first replays every item
Partition Keystringno-Follow only this one partition
Poll Intervaldurationno5sHow long to wait before asking again once the feed is idle (1s–5m)
Max Items Per Pollintegerno100Most items fetched in one read, 1–1,000
Timeoutdurationno-Cap on how long one follow may run. Leave empty
What the change feed does and does not report

The change feed reports the latest version of each item that was created or updated:

  • Deletes are not reported. An item deleted before the feed is read never appears at all. To act on deletions, mark items deleted (a soft-delete field, with a TTL to remove them later) instead of deleting them.
  • Rapid updates collapse. An item written twice between two reads arrives once, with its latest content.

The feed's position is kept in memory. After MaestroHub restarts, the trigger starts again from its Start From setting: now skips changes made while it was down, and beginning replays the container. For this reason the connection is exclusive — only one replica runs it, since a second would read the same feed and fire each change twice.

Using Parameters​

Every field marked templatable accepts ((parameterName)) placeholders, filled from upstream node output at run time:

WhereExample
Item ID / Partition Key((deviceId))
Item{"id": "((id))", "deviceId": "((deviceId))", "celsius": ((value))}
Query Parameters{"@device": "((deviceId))"}
Operations[{"op": "set", "path": "/status", "value": "((status))"}]
Parameter configuration for an Azure Cosmos DB function

Parameters detected from ((placeholders)) in a Write Item function

A templated value is only known at run time, so the form defers its checks on fields that contain a placeholder; the same checks run on the resolved value when the function executes.

Parameter Availability

Every function type except Change Feed accepts parameters. The change feed starts a pipeline, so there is no upstream node to fill them from.

Pipeline Integration​

Each execution function appears as its own node under Database in the node library, and the change feed as an Azure Cosmos DB Trigger:

Azure Cosmos DB Write Item node

An Azure Cosmos DB Write Item node on the canvas

Common patterns:

  • Store readings — an MQTT or OPC UA trigger followed by Write Item
  • Enrich — a Read Item or Read Many Items lookup before a decision node
  • Count and stamp — Patch Item with incr and set in one atomic call
  • Keep two items in step — a Transactional Batch that writes a reading and updates its device's summary together

Store-and-forward and replay safety​

The four write operations can be buffered through store-and-forward. What a replay does differs per operation:

OperationOn replay
Write Item, upsert or replaceWrites the same body again — same end state
Write Item, createFails with 409 — the first write stands
Patch ItemApplies again — an incr adds twice
Delete ItemAnswers deleted false — same end state
Transactional BatchApplies whole or fails whole, like the first attempt

Patch Item is declared non-idempotent because of incr. If a patch only sets, replaces or removes fields, a replay converges on the same state.

Common Use Cases​

Machine state store​

Keep one item per machine, partitioned on /deviceId. An OPC UA trigger patches the machine's item with set for its status and incr for its cycle counter, and dashboards read it with Read Item.

Time-series readings​

Write each reading as its own item with Write Item, partitioned by device. A scheduled pipeline runs Query Items with SELECT * FROM c WHERE c._ts > @since scoped to one device's partition to hand the latest readings to an analytics system.

Order intake​

Follow an orders container with the change feed. Each new order runs a fulfilment pipeline, which writes the shipment and updates the order in one Transactional Batch.

Replicating to another system​

Follow a container from the beginning once to copy it into a warehouse, then keep following to forward every later change.

Troubleshooting​

SymptomCauseFix
database "…" not found in this accountThe Database field names a database that does not existCheck the name; database names are case-sensitive
401 Unauthorized on connectA wrong account keyCopy the key again from the account's Keys blade
401 Unauthorized on writes onlyA read-only keyUse the primary or secondary read-write key
403 Forbidden with Entra IDThe identity has no data-plane roleAssign Cosmos DB Built-in Data Contributor with az cosmosdb sql role assignment create
container "…" not foundA typo in the Container fieldContainer names are case-sensitive
item has no value at /…The item lacks the container's partition key fieldAdd the field to the item
409 ConflictCreate found the ID takenUse upsert, or a new ID
412 Precondition FailedThe ETag changed, or the patch condition did not matchRead the item again and retry with its new _etag
413 Request Entity Too LargeAn item over 2 MBSplit the item
429 Too Many RequestsThe container is over its RU/sRaise the throughput, or spread the load
A query error mentioning cross partitionORDER BY, TOP, DISTINCT, GROUP BY or an aggregate across partitionsSet Partition Key
Read Item finds nothing for a numeric key42 is the string keyWrite the key as [42]
The trigger misses deletesThe change feed does not report deletesSoft-delete with a field and a TTL