For AI agents: a documentation index is available at https://www.mongodb.com/docs/llms.txt — markdown versions of all pages are available by appending .md to any URL path.
See how MongoDB 9.0 delivers up to 2x higher throughput.
MongoDB Branding Shape
Register now >
Docs Menu

$throttle Aggregation Stage (Stream Processing)

$throttle

The $throttle stage limits the rate at which a stream processor passes data to the stages that follow it. Use $throttle to protect downstream systems from bursts of traffic and to stay within the rate limits that those systems enforce.

Bursts can occur when a stream processor takes an initial snapshot of a large data set, recovers after an outage, or reprocesses historical data. Without a rate limit, these bursts can overwhelm the target system, force it to scale up, or exceed a third-party API quota.

A $throttle stage has the following prototype form:

{
$throttle: {
bytesPerSec: <integer>,
messagesPerSec: <integer>
}
}

The $throttle stage takes a document with the following fields:

Field
Type
Necessity
Description

bytesPerSec

integer

Optional

Maximum number of bytes per second that the stage passes to downstream stages. Must be a positive integer.

messagesPerSec

integer

Optional

Maximum number of messages per second that the stage passes to downstream stages. Each document counts as one message. Must be a positive integer.

You must specify at least one of bytesPerSec or messagesPerSec. If you specify neither, Atlas Stream Processing returns the following error:

$throttle requires at least one of 'bytesPerSec' or 'messagesPerSec'

Atlas Stream Processing tracks a separate rate limit for each field that you set. Each limit regains capacity continuously at the rate that you configure, and Atlas Stream Processing tracks each limit independently of the others.

Messages can flow past a $throttle stage only when every configured limit has available capacity for them. If you set both fields, both limits must have capacity. When more than one limit is exhausted, the stage waits for the time that the most exhausted limit requires. All limits regain capacity during that wait, so a limit with a smaller shortfall recovers before the wait ends.

The stage doesn't wait for a full second before it passes more data. It processes the next messages as soon as every configured limit has enough capacity for them.

If a single document is larger than the bytesPerSec limit, the stage doesn't hold the document until enough capacity accumulates. Instead, the stage passes the document to the next stage and then waits one second before it passes more data.

You can use more than one $throttle stage in a single pipeline. Each stage limits only the data that flows through it, which lets you apply different limits to different parts of your pipeline.

You can use $throttle only in the main pipeline. You can't nest $throttle inside another stage.

Throttling slows the flow of data through your pipeline, which creates backpressure on the stages that precede the $throttle stage. If your source produces data faster than your throttle limit allows, the processor falls behind the source and lag grows over time.

To avoid unbounded lag, set limits that match the sustained throughput of your source rather than only its peak. Monitor stats.changeStreamTimeDifferenceSecs and the throttle statistics described in Monitoring to confirm that your processor keeps up with its source.

Atlas Stream Processing reports the following statistics for stream processors that use $throttle:

Statistic
Description

throttle.throttledTimeMs

Cumulative time in milliseconds that the processor spent actively throttling.

throttle.throttleEvents

Number of times that the stage delayed messages to stay within the configured rate limits.

Atlas Stream Processing reports both statistics at the processor level and for each stage. When a pipeline contains more than one $throttle stage, the processor-level values are the sum of the per-stage values. The per-stage statistics are stats.operatorStats.throttle.throttledTimeMs and stats.operatorStats.throttle.throttleEvents.

To learn more about stream processor statistics, see Atlas Stream Processing Monitoring and Metrics.

The following pipeline limits the data that a stream processor writes to a Kafka topic to 5 MiB and 50 messages per second:

{
$throttle: {
bytesPerSec: 5242880,
messagesPerSec: 50
}
},
{
$emit: {
connectionName: "ordersTopic"
}
}

Consider a batch of 20 documents totaling 2 MiB. Recent traffic has used some of the capacity of both limits:

Limit
Capacity needed
Capacity available
Shortfall
Recovery rate
Required wait

messagesPerSec

20

15

5

50/sec

100 ms

bytesPerSec

2,097,152

1,048,576

1,048,576

5,242,880/sec

200 ms

The stage can't proceed until both limits have enough capacity, so it waits 200 ms, the longer of the two wait times. Because both limits regain capacity during that wait, the messagesPerSec limit, which needed only 100 ms, has already recovered when the wait ends.