Skip to main content

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​

FieldTypeRequiredDefaultDescription
keyExpressionStringYes—FEEL expression evaluated on each record. Records with the same result share one limit.
numToAllowIntegerNo1Records allowed per key in each period. Between 1 and 2147483647.
periodSecondsLongNo30Length of each period in seconds. Between 1 and 31536000 (365 days).
cacheSizeLimitIntegerNo50000Maximum 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 null or 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 timeInputOutput
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: []
  • filter: Drop records that don't match a condition
  • eval: Compute a key field before throttling