Atlas Stream Processing は、階層に従ってストリーム プロセッサごとにリソースを割り当てます。リソース割り当てとコストが固定されることで予測可能性が得られ、システム設計プロセスが簡素化されます。このガイドを使用して、配置を計画する際に、どの階層がストリーム処理ワークロードに最も適しているかを理解してください。
リソース割り当て
各階層は、処理能力、メモリ、帯域幅、並列処理、およびApache Kafkaソースを持つプロセッサの場合、パーティションの固定の割り当てを提供します。
階層 | vCPU | RAM(GB) | 帯域幅 (MB) | 最大並列処理 | ソースKafkaパーティションの制限 | 最初の同期のコレクション制限 |
|---|---|---|---|---|---|---|
SP 2 | 0.25 | 0.5 | 50 | 1 | 32 | 1 |
SP 5 | 0.5 | 1 | 125 | 2 | 64 | 1 |
SP 10 | 1 | 2 | 200 | 8 | 無制限 | 5 |
SP 30 | 2 | 8 | 750 | 16 | 無制限 | 10 |
SP 50 | 8 | 32 | 2500 | 64 | 無制限 | 50 |
並列処理
並列処理は、ストリーム プロセッサがデータの読み取り、豊富、書込み (write) に使用できるスレッドまたは同時要求の数を決定します。個々のパイプラインステージで並列処理を構成しますが、Atlas Stream Processing は、前の表に示されているプロセッサの階層の最大値に対してプロセッサ全体にそれを強制します。
並列処理をサポートするステージ
次のステージは parallelism 値を受け入れます。各 のデフォルトは 1 です。
ステージ | 値が大きい場合の影響 |
|---|---|
| |
| |
Atlas Stream Processing が書込み (write) 操作を分散するスレッド数を増やします。そのため、ストリーム プロセッサとそれが書込むクラスターの両方でより多くの計算リソースを使用する必要があります。 | |
Sink 演算子が使用する内部書込みスレッドの数を増やし、それらのスレッド全体に書込み操作を分散します。 | |
外部関数に対する並列リクエストの最大数を増やすため、より多くの計算リソースが必要になります。 |
実行時に並列処理の仕組み
ストリーム プロセッサは、次のシーケンスを通じてデータを移動します。
Source -> Buffer -> Transform -> Buffer -> Sink
Apache Kafkaソースの場合、Atlas Stream Processing はソース パーティションごとに 1 つのコンシューマー スレッドを開始するため、プロセッサは一度に多くのパーティションから読み取ることができます。パイプラインの中央にある変換ステージは引き続き単一のスレッドで実行され、バッファされたデータをバッチで処理します。
1 より大きい parallelism 値を持つステージがバッチするを受け取ると、そのバッチするに並列スレッドまたはリクエストが実行されます。その後、バッチするは次のステージに移動します。
ドキュメントの分散と順序
ステージの parallelism 値が 1 より大きい場合、Atlas Stream Processing は次のようにドキュメントをスレッドに分散します。
partitionByを指定しない場合、Atlas Stream Processing はラウンドロギングの順序でドキュメントをスレッドに割り当てます。partitionByを指定すると、Atlas Stream Processing は同じpartitionBy値を持つすべてのドキュメントを同じスレッドに送信し、ドキュメントを順番に処理します。
$merge onは例外です。これは、 句内のフィールドをパーティションキーとして使用します。
累積並列処理
各ストリーム プロセッサには、その階層によって決定される最大累積並列処理値があります。ストリーム プロセッサの累積並列処理は、次のように計算されます。
parallelism total - parallelized stages
parallelism total は、$source、$lookup、$merge、$emit、および $externalFunction ステージにおける 1 を超えるすべての parallelism 値の合計であり、parallelized stages はこれらのステージの数で parallelism の値が 1 より大きい。
例、$source ステージが 4 の parallelism 値を設定し、$lookup ステージでは parallelism 値が設定されていない(デフォルトは 1)、かつ $merge ステージが parallelism に設定されている場合: 2 の } 値がある場合は 2 つの parallelized stages があり、ストリーム プロセッサの累積並列処理は (4 + 2) - 2 として計算されます。
ストリーム プロセッサがその階層の最大累積並列処理を超える場合、Atlas Stream Processing はエラーをスローし、目的のレベルの並列処理に必要な最小プロセッサ階層について提案します。エラーを解決するには、プロセッサをより高い階層に増やすアップするか、ステージの並列処理値を低くする必要があります。詳しくは、Stream Processingをご覧ください。
Kafkaソース パーティション
Apache Kafkaから読み取るプロセッサも、階層のソース パーティション制限によって制約されます。 SP2 階層ではプロセッサが32 ソース パーティションに制限され、SP5 階層ではプロセッサが に制限されます。64 SP10 階層 以上ではパーティション制限はありません。
階層のパーティション制限を超えるプロセッサは失敗するため、追加のパーティションをサポートするには増やすアップする必要があります。プロセッサの実行中にトピックがパーティションを増やす可能性があるため、予想されるパーティションの増加に対応できるヘッドルームのある階層を選択します。 Kafkaソースの動作の詳細については、「 Atlas Stream Processing の制限 」を参照してください。
構成された並列処理の監視
stats.addedParallelism構成した並列処理とプロセッサの階層で使用可能な並列処理を比較するには、プロセッサの統計情報が返す フィールドを使用します。 Atlas Stream Processing では、少なくとも 1 つのステージでparallelism 値が1 より大きい場合にのみこのフィールドが返されます。
ワークロードの選択
各階層の異なるリソース割り当ては、プロジェクトの異なるステージと規模に適しています。
階層 | ユースケース |
|---|---|
SP 2 | 開発、 試用配置 限られたリソース要件で基本的なワークロードをサポートできる、最もコストの低いオプション。 |
SP 5 | 開発、基本的な本番環境の配置 より複雑な計算を使用する場合でも、 低スループットの本番タスクに適した低コストのオプション。SP5 プロセッサは、基本的なフィルタリング、プロジェクション、変更ストリーム処理をサポートできます。 |
SP 10 | メインストリームの本番環境への配置 本番環境ワークロードのベースライン。SP10 以上は、より高いレベルの並列処理、無制限のKafkaパーティショニング、または lookup や join などのデータ豊富操作を必要とするパイプラインを対象としています。 |
SP 30 | 複雑な本番環境の配置 メモリ集中型のステートメント操作用に設計された高性能オプション。 SP30 は、長時間のウィンドウ、複数の検索、および増やすでデータを増やすために大容量のRAMバッファを必要とするステージを使用するパイプラインをサポートします。 |
SP 50 | エンタープライズスケールの本番環境 高スループット ストリームと広範な変換ロジック向けに設計された最もパフォーマンスの高いオプション。SP50 プロセッサは、大規模な並列処理を必要とする操作やコンピューティング集中型のワークフローに適しています。 |
Considerations
適切な階層を選択する際には、次の要素を考慮する必要があります。
初期化サージ
ストリーム プロセッサは、通常の操作中よりも初期実行中に多くのリソースを必要とする場合があります。例、大規模な Atlasコレクションに対して initialSync を実行するプロセッサは、同期中に重い I/O と計算をサポートする必要があります。
このような需要の上昇を対応するには、一時的に高い階層を選択し、同期が完了し、プロセッサが新しい変更ストリームイベントのみを消費するように移行したときにプロセッサを増やすします。
パイプライン ロジック
集計パイプラインライン ロジックは、CPU とRAM消費のプライマリ ドライバーです。
- Windows: 長時間続くウィンドウでは、移動中に保持するRAMがより多く消費されます
- ドキュメント.
- カスタムロジック: Javascript
$functionステージまたは複雑なグループ化 - ロジックにより、各メッセージの計算要件が増加します。
- カスタムロジック: Javascript
- 複合複雑さ: 追加のステートメントまたは計算が複雑なステージ
- リソース需要の潜在的な変化がより大きくなります。過剰なキャパシティーを維持することで、消費量が急増しても一貫したスループットが確保されます。
インフラストラクチャ
ネットワークまたはストレージに接続する点に、ストリーム プロセッサのオーバーヘッドが増加します。
- ソースまたはシンク密度: からの読み取りまたは並列化されたへの書込み
- ソースまたは シンク(Apache Kafkaトピックなど)では、I/O 要件が増加します。
- データ強化:
$lookupステージと$httpsステージおよび 操作 - Atlas コレクションに対してストリーム内のデータを増やすには、 ネットワーク帯域幅 と接続プーリング が必要です。
- データ強化:
- 調整: 多くのソースをオーケストレーションする複雑な配置で
- と シンクの場合、ストリーム プロセッサはこれらの各ノード間でデータ フローをルーティングするハブとして機能できます。このようなプロセッサでは、SP30 および SP50 階層のより高いスループットのメリットが得られます。
逆に、高スループットのストリーム処理ワークロードでは、接続されたリソースの需要が増加する可能性があります。
- Atlas への影響: 大ボリューム、ストリームからの並列 I/O
プロセッサは、ソースまたはシンク Atlas クラスターの読み取りまたは書込みキャパシティーを超えることができます。これによりプロセッサのレイテンシが増加するだけでなく、それらのクラスターに依存する他のワークロードがボトルネックになる可能性もあります。
システム全体のパフォーマンスを確保するには、Atlas クラスターを、交流するプロセッサに比例して増やす。
パフォーマンスとレイテンシ
処理ロジックが最小限の場合でも、パフォーマンス目的では上位階層のプロセッサが必要になる場合があります。
- 高スループット: 上位階層プロセッサほど、ストリームのサポートが向上します
- 高いレートで イベントを生成します。
低レイテンシ SLA: 上位階層のプロセッサが提供する高並列処理により、速度が重要な場合にイベントがキューに蓄積されなくなります。特に、SP50 プロセッサは SP30 プロセッサの 4 倍のスレッドを提供します。
- データ暗号化とキャッシュ:
$cachedLookupを使用してデータを増やす場合 - ストリーム間で大規模な静的参照データまたは低速変化参照データを含む場合は、キャッシュに必要なRAM を提供するために上位階層のプロセッサが優先されます。
- データ暗号化とキャッシュ:
- 複雑なシンク: 特定のシンクはよりコストのかかる
- 変換、トランザクション、ファイル マネジメントのオーバーヘッドが含まれます。これらのシンクと交流するプロセッサの場合、上位の階層はコンシステントなパフォーマンスとレイテンシを確保するのに役立ちます。
スケーリング
Atlas Stream Processing のスケーリングは垂直型です。プロセッサを手動で増やすことも、オートスケーリングを有効にして Atlas Stream Processing を使用して階層を調整することもできます。
プロセッサを手動で増やすには、プロセッサを停止して新しい階層を選択して、再起動します。 Atlas Stream Processing チェックポイント により、移行中にデータが失われることはありません。プロセッサのパフォーマンスを定期的に監視し、このガイドに記載されている要素に基づいて階層を調整します。
Atlas Stream Processing がリソースの使用状況に応じて階層を自動的に調整できるようにするには、垂直オートスケーリングを有効にします。このガイドの要素を使用して、プロセッサが増やすできる範囲を境界とするminTier とmaxTier を選択します。