定義
$throttle ステージは、ストリーム プロセッサがそれに続くステージにデータを渡すレートを制限します。 $throttle を使用して、下流のシステムをトラフィックのバーストから保護し、それらのシステムが強制するレート制限内にします。
ストリーム プロセッサが大規模なデータセットの 最初のスナップショット を取得し、停止時後に回復 するか、履歴データを再処理する場合にバーストが発生する可能性があります。レート制限がない場合、これらのバーストはターゲット システムを超え、増やすアップを強制したり、サードパーティのAPI割り当てを超えたりする可能性があります。
構文
$throttle ステージのプロトタイプ形式は次のとおりです。
{ $throttle: { bytesPerSec: <integer>, messagesPerSec: <integer> } }
$throttleステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| integer | 任意 | ステージが下流の ステージに渡す 1 秒あたりの最大バイト数。正の整数である必要があります。 |
| integer | 任意 | ステージが下流の ステージに渡す 1 秒あたりのメッセージの最大数。各ドキュメントは1 つのメッセージとしてカウントされます。正の整数である必要があります。 |
bytesPerSec または messagesPerSec の少なくとも 1 つを指定する必要があります。どちらも指定しない場合、Atlas Stream Processing は次のエラーを返します。
$throttle requires at least one of 'bytesPerSec' or 'messagesPerSec'
動作
Atlas Stream Processing は、設定した各フィールドのレート制限を個別に追跡します。各制限はユーザーが設定したレートでキャパシティーを継続的に回復し、Atlas Stream Processing は各制限を他の制限と独立して追跡します。
メッセージは、設定されたすべての制限で使用可能なキャパシティーがある場合にのみ、 $throttle ステージを過去にフローすることができます。両方のフィールドを設定する場合、両方の制限のキャパシティーは である必要があります。複数の制限が使い果たされた場合、 ステージは最も消費された制限に必要な時間のみ待機します。すべての制限は待機中にキャパシティーを回復するため、待機が終了する前に短縮値が小さい制限は回復します。
ステージは、さらにデータを渡す前に完全な秒を待機しません。設定されたすべてのメッセージのキャパシティーが十分になると、すぐに次のメッセージを処理します。
1 つのドキュメントが bytesPerSec の制限を超える場合、十分なキャパシティーが蓄積されるまで、 ステージはドキュメントを保持しません。代わりに、 ステージはドキュメントを次のステージに渡し、1 秒間待ってからさらにデータを渡します。
単一のパイプラインで複数の $throttle ステージを使用できます。各ステージはそれを通過するデータのみを制限するため、パイプラインのさまざまな部分に異なる制限を適用できます。
$throttle は メインパイプラインでのみ使用できます。別のステージ内に $throttle をネストすることはできません。
バックプレッシャー
スロットリングにより、パイプラインを介するデータの流れが遅くなり、$throttle ステージに先行するステージにバックプレッシャーが生じます。ソースがスロットル制限の許可よりも速くデータを生成する場合、プロセッサはソースより遅れ、ラグは時間の経過とともに大きくなります。
無制限の遅延を回避するには、ソースのピークのみではなく、ソースの持続的なスループットに一致する制限を設定します。 と モニタリングstats.changeStreamTimeDifferenceSecs で説明されているスロットルの統計情報をモニターして、プロセッサがソースに対応していることを確認します。
モニタリング
Atlas Stream Processing は、$throttle を使用するストリーム プロセッサについて次の統計を報告します。
統計値 | 説明 |
|---|---|
| プロセッサがアクティブなスロットリングに費やした累計時間(ミリ秒単位)。 |
| 設定されたレート制限内に収まるように ステージがメッセージを遅延させた回数。 |
Atlas Stream Processing では、プロセッサ レベルの統計情報と各ステージの統計情報の両方が報告されます。パイプラインに複数の $throttle ステージが含まれている場合、プロセッサレベルの値はステージごとの 値の合計です。ステージごとの統計は stats.operatorStats.throttle.throttledTimeMs と stats.operatorStats.throttle.throttleEvents です。
ストリーム プロセッサの統計の詳細については、「 Atlas Stream Processing の監視とメトリクス 」を参照してください。
例
次のパイプラインでは、ストリーム プロセッサがKafkaトピックに書き込むデータを、1 秒あたり 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 ミリ秒 |
ステージは両方の制限が十分なキャパシティーになるまで続行できないため、2 つの待機時間のうち長い方の 200 ミリ秒待機します。両方の制限は待機中にキャパシティーを回復するため、100 ミリ秒のみを必要とした messagesPerSec の制限は、待機が終了するとすでに回復しています。