定义
$throttle 阶段限制了流处理器将数据传递到其后续阶段的速率。使用 $throttle 可保护下游系统免受流量突发的影响,并保持在这些系统实施的速率限制之内。
当流处理器拍摄大型数据集的初始快照、中断后恢复或重新处理历史数据时,可能会发生突发。如果没有速率限制,这些突发可能会使目标系统不堪重负,迫使其扩展,或超过第三方API配额。
语法
$throttle 阶段具有以下原型形式:
{ $throttle: { bytesPerSec: <integer>, messagesPerSec: <integer> } }
$throttle 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 整型 | Optional | 该阶段每秒传递给下游阶段的最大字节数。必须为正整数。 |
| 整型 | Optional | 该阶段每秒传递给下游阶段的最大消息数。每个文档算作一条消息。必须为正整数。 |
您必须至少指定 bytesPerSec 或 messagesPerSec 之一。如果两者均未指定, Atlas Stream Processing将返回以下错误:
$throttle requires at least one of 'bytesPerSec' or 'messagesPerSec'
行为
Atlas Stream Processing跟踪您设立的每个字段的单独速率限制。每个限制都会按照您配置的速率持续重新获得容量,并且Atlas Stream Processing会独立跟踪每个限制。
仅当每个配置的限制都有可用容量时,消息才能流经 $throttle 阶段。如果同时设立这两个字段,则两个限制都必须具有容量。当多个限制用尽时,该阶段会等待最用尽的限制所需的时间。在此等待期间,所有限制都会重新获得容量,因此缺口较小的限制会在等待结束之前恢复。
该阶段不会等待一整秒,就会传递更多数据。一旦每个配置的限制都有足够的容量来处理接下来的消息,它就会处理这些消息。
如果单个文档大于 bytesPerSec 限制,则阶段不会容纳该文档,直到积累了足够的容量。相反,该阶段会将文档传递到下一阶段,然后等待一秒钟,然后再传递更多数据。
您可以在单个管道中使用多个 $throttle 阶段。每个阶段仅限制流经它的数据,这样您就可以对管道的不同部分应用不同的限制。
您只能在主管管道中使用 $throttle。您不能将 $throttle 嵌套在另一个阶段内。
背压
限制会减慢通过管道的数据流,从而在 $throttle 阶段之前的阶段上产生反压。如果源生成数据的速度超过限制允许的速度,则处理器会落后于源,并且延迟会随着时间的推移而增加。
为避免无限制的延迟,设立与源的持续吞吐量(而不仅仅是其峰值吞吐量)相匹配的限制。监控 stats.changeStreamTimeDifferenceSecs 和监控中描述的限制统计信息,以确认您的处理器与其源保持同步。
监控
Atlas Stream Processing报告使用 $throttle 的流处理器的以下统计信息:
统计信息 | 说明 |
|---|---|
| 处理器主动限制所花费的累计时间(以毫秒为单位)。 |
| 该阶段为保持在配置的速率限制内而延迟消息的次数。 |
Atlas Stream Processing报告处理器级别和每个阶段的统计信息。当管道包含多个 $throttle 阶段时,处理器级别的值是每阶段值的总和。每个阶段的统计信息为 stats.operatorStats.throttle.throttledTimeMs 和 stats.operatorStats.throttle.throttleEvents。
示例
以下管道流处理器写入Kafka主题的数据限制为每秒 5 MiB 和 50 条消息:
{ $throttle: { bytesPerSec: 5242880, messagesPerSec: 50 } }, { $emit: { connectionName: "ordersTopic" } }
以批处理总 2 MiB 的 20 文档为例。最近的流量已使用这两个限制的部分容量:
Limit | 所需容量 | 可用容量 | 不足 | 恢复率 | 所需等待 |
|---|---|---|---|---|---|
| 20 | 15 | 5 | 50/秒 | 100 毫秒 |
| 2,097,152 | 1,048,576 | 1,048,576 | 5,242,880/秒 | 200 毫秒 |
在两个限制都有足够的容量之前,该阶段无法继续,因此它会等待 200 毫秒,即两个等待时间中较长的一个。由于两个限制在等待期间都重新获得了容量,因此当等待结束时,messagesPerSec 限制已经恢复,该限制只需要 100 毫秒。