Throttle Node
The throttle node is a per-key rate limiter. It groups records by a key computed with a FEEL expression and lets at
most numToAllow records per key through in each period of periodSeconds. Records over the limit are dropped. When a
key sends records again after a period with drops, the first record let through carries a throttledCount field with
the number of records dropped in that period.
For more details of using FEEL, please refer to FEEL Reference
Key Features
- Per-key limits: Each distinct key value has its own counter, so a noisy source can't use up the limit of a quiet one
- Fixed periods: A period starts at the first record of a key after the previous period has ended and lasts
periodSeconds - Drop accounting: The first record let through after a period with drops carries
throttledCount, so downstream consumers know how many records were suppressed - Fail-open: Records whose key can't be computed pass through unthrottled instead of being dropped
- Bounded memory: The number of tracked keys is capped by
cacheSizeLimit
Config Object
| Field | Type | Required | Default | Description |
|---|---|---|---|---|
keyExpression | String | Yes | — | FEEL expression evaluated on each record. Records with the same result share one limit. |
numToAllow | Integer | No | 1 | Records allowed per key in each period. Between 1 and 2147483647. |
periodSeconds | Long | No | 30 | Length of each period in seconds. Between 1 and 31536000 (365 days). |
cacheSizeLimit | Integer | No | 50000 | Maximum number of keys to track. Between 1 and 2147483647. See Key Tracking Limit. |
Numeric fields must be whole numbers: 3.0, "3" and out-of-range values are rejected. Unknown fields are rejected
too, so a misspelled field name fails validation instead of being ignored.
How Throttling Works
Periods
For each key, the first record opens a period of periodSeconds. The first numToAllow records of the key in that
period pass; the rest are dropped. The first record at or after the end of the period opens a new period and the
counter resets.
The period is anchored at the first record of the period, not at the last allowed one. With numToAllow: 2 and
periodSeconds: 60, if records of the same key arrive at 0s and 40s, both pass, a record at 50s is dropped, and a
record at 61s passes because the period [0s, 60s) has ended.
Time is the node's processing time (the wall clock when the record reaches the node), not a timestamp inside the record.
The throttledCount Field
When a key opens a new period and records were dropped in its previous period, the first record of the new period is
passed with an extra top-level field throttledCount set to the number of records dropped. If the record already has
a throttledCount field, it is overwritten. No field is added when nothing was dropped.
The count is attached only when the key sends another record. If a key stops sending, its last drop count is not reported on any record.
Records Without a Usable Key
A record passes through unthrottled, without throttledCount, when its key expression:
- returns
nullor an empty string - returns an object or array
- fails to evaluate
Keys are compared by their string form, with whole numbers normalized. The number 5, the number 5.0 and the
string "5" are all the same key, and so are true and "true".
Key Tracking Limit
Per-key state is kept in memory. Every 10,000 records with a usable key, if more than cacheSizeLimit keys are
tracked, the node first drops keys with no records for two periods, then the least recently seen keys until it is back
at the limit. The number of tracked keys can therefore exceed cacheSizeLimit for a while. A key that is dropped from tracking
starts a new period on its next record, and any drop count it had is not reported.
Scope of State
State lives in the memory of each running pipeline instance. It is not shared across instances and is lost when the pipeline restarts. If a key's records are spread across several instances, each instance applies its own limit.
Metrics
In addition to the standard input_event_count, output_event_count and error_event_count counters, the node emits
throttle_dropped_count for each dropped record. error_event_count counts records whose key expression failed to
evaluate; those records still pass through.
Examples
Composite Key
Combine fields into one string to throttle each host and error code pair separately. The key must be a single
value: an expression returning a dict or an array is not a usable key, so those records pass through unthrottled.
{
"keyExpression": "concat($.host, \"|\", $.error_code)",
"numToAllow": 1,
"periodSeconds": 300
}
concat converts each value to a string and joins them. A missing field becomes an empty string, so a record
without error_code is still throttled, under a key such as "web-1|". Whole-number normalization applies only to
the final key, so an error_code of 500 and of 500.0 give the different keys "web-1|500" and "web-1|500.0".
Sample Input and Output
With keyExpression: "$.host", numToAllow: 2 and periodSeconds: 60:
| Arrival time | Input | Output |
|---|---|---|
| 0s | {"host": "a", "seq": 1} | {"host": "a", "seq": 1} |
| 40s | {"host": "a", "seq": 2} | {"host": "a", "seq": 2} |
| 50s | {"host": "a", "seq": 3} | dropped |
| 55s | {"host": "b", "seq": 4} | {"host": "b", "seq": 4} |
| 61s | {"host": "a", "seq": 5} | {"host": "a", "seq": 5, "throttledCount": 1} |
| 61s | {"host": "a", "seq": 6} | {"host": "a", "seq": 6} |
| 61s | {"host": "a", "seq": 7} | dropped |
| 121s | {"host": "a", "seq": 8} | {"host": "a", "seq": 8, "throttledCount": 1} |
DAG Definition
jobContext:
metricTags: {}
otherProperties: {}
dag:
- id: "source"
commandName: "kafkasource"
config:
broker: "localhost:9092"
topic: "alerts"
groupId: "alert-throttle"
encodingType: "JSON_OBJECT"
outputs:
- "throttle"
- id: "throttle"
commandName: "throttle"
config:
keyExpression: "$.host"
numToAllow: 1
periodSeconds: 60
outputs:
- "sink"
- id: "sink"
commandName: "stdout"
config:
encodingType: "JSON_OBJECT"
outputs: []