データ接続と統合パイプラインの構築ストリーミングパイプラインストリーミングのステートフル変換

注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。

ストリーミングのステートフル変換

Foundry パイプラインは、複雑なデータ加工を高速なストリーミングで実行できるステートフルなデータストリーミング変換を提供します。

状態とは?

データストリーミングでは、各出力行が以前に処理された行の情報に依存する可能性がある場合、そのデータ変換はステートフルです。行の処理をまたいで保持される情報を状態と呼びます。各行の処理で、状態にアクセスして変更できます。

ステートフルなデータ加工の例として、合計を求める集計データ加工があります。次のサンプルテーブルでは、合計を求める集計データ加工を使用して、値が hot であるセンサー読み取り値の件数を日ごとに計算できます。ストリーミングで受信した時刻が新しい行ほど、テーブルの上部に表示されています。

日センサー読み取り値タイムスタンプ
Mondayhot3
Mondaycold2
Mondayhot1

日ごとに hot の読み取り値の累計件数をライブで計算する、ステートフルな合計集計データ加工の出力は、次のようになります。

日センサー読み取り値タイムスタンプ状態
Mondayhot32
Mondaycold21
Mondayhot11

ステートフルなデータ加工は、行全体を含め、シリアル化可能な任意のデータ型の状態を保存できるため、複雑な動作を実現できます。Pipeline Builder の状態を扱うデータ加工はあらかじめ用意されており、状態のデータ型とその変化を自動的に処理します。

ユーザーが状態の内容に常にアクセスできるとは限りません。状態は、バックエンドでの処理を実現するために使用できます。たとえば、データパイプラインの Flink ジョブでストリーム入力ソースを読み込むと、パーティションごとに最後のオフセットの状態が保存されます。これにより、障害や再起動の後に、最後に成功した時点からジョブを復旧できます。

データストリーミングで状態が強力なのはなぜですか?

Foundry のデータストリーミングは Flink アーキテクチャを使用して、低遅延のデータパイプライン処理を提供します。各行は計算され、処理後すぐに下流の次の操作に渡されます。バッチ変換とは異なり、ストリーミング変換はデータ全体を参照せずに、1行ずつ次の出力を決定します。

状態を持たない(ステートレスとも呼ばれる)ストリーミング変換では、このアーキテクチャにより、データ加工のロジックが一度に依存できるのは1行だけになります。たとえば、ステートレスなデータ加工では、Integer 列に常に5を加算できます。

これに対して、ステートフルなストリーミングデータ加工は、行が到着するたびに1行ずつ即座に処理しながら、以前の行について永続化されたデータにアクセスできます。

Foundry のステートフルストリーミングでは、Exactly-once 保証のオプションが用意されており、これがデフォルトのパイプライン設定です。このオプションを選択すると、状態を変更する行によって、キーごとの順序で、厳密に1回だけ状態が変更されることが保証されます。これにより、正確で複雑なデータストリーミング動作を実現できます。

たとえば、合計を求める場合は、ストリームが再起動したり障害が発生したりしても、合計は常に正確になります。並べ替えを使用する場合は、ジョブがライブではない部分的な並べ替え済み出力を生成した後、並べ替えの途中でジョブを再起動しても、入力ごとに必ず1つだけ出力が生成されます。

キー付き状態

Pipeline Builder のステートフルなストリーミング変換はすべてキー付き状態を使用し、ユーザーがパーティションキー列を指定する必要があります。ステートフルなデータ加工では、キー列の値が異なる行は別々に処理されます。これにより、バックエンドで処理を並列化し、大量のデータに対応できます。

たとえば、Day 列をパーティションキーとして使用し、日ごとに hot の読み取り値の累計件数をライブで計算する、ステートフルな合計集計の例を考えてみます。

日センサー読み取り値タイムスタンプ状態
Tuesdayhot51
Mondayhot42
Tuesdaycold30
Mondaycold21
Mondayhot11

Day の値が Monday の行と Tuesday の行では、状態が独立して計算されている点に注目してください。キーの値が Tuesday の行が到着しても、Monday に対して保存される内容には影響しません。また、Day 列の値が Wednesday などの異なる値である行がさらに到着しても、これらのキーの状態は影響を受けません。

レコードの分散が非効率になるキーは、不要な負荷の増加やスループットの制限につながる可能性があるため、キーは慎重に選択してください。ストリーミングキーのベストプラクティスを参照してください。

イベント時刻とウォーターマーク

ステートフルなストリーミングデータ加工は、時間に関する情報に依存することがよくあります。ストリームは継続的に流れ、いつでも新しいリアルタイムの行を受信する可能性があるため、時間的に近い行をグループ化することが適切な場合が多くあります。たとえば、外部キャッシュ結合データ加工は、結合列の値が一致し、行のタイムスタンプの差が有効期限の範囲内にある場合にのみ、2つの入力ストリームの行を結合します。Pipeline Builder のストリーミングは Flink のイベント時刻を使用して、リプレイ時に同じ、または非常によく似た出力となる、決定論的に近い方法でステートフルな変換を実現します。

Pipeline Builder で時間に基づく操作を行うステートフルなデータ加工には、上流にタイムスタンプとウォーターマークを割り当てるデータ加工が必要です。パイプライングラフにこのデータ加工がない場合は、検証エラーが発生します。タイムスタンプとウォーターマークの割り当てでは、各行に「イベント時刻」を割り当てます。通常は、その行に含まれるタイムスタンプ列を使用します。ウォーターマークは、各データ加工操作における「現在時刻」のほぼ決定論的な指標です。単調に増加する値であり、その演算子への入力行で確認されたイベント時刻の最大値の少し後を追うように進みます。たとえば、外部キャッシュ結合がキャッシュ内のエントリーの有効期限が切れたかどうかを判定する際には、ウォーターマークが有効期限の時刻以上かどうかを確認します。これは、そのイベント時刻以降の入力行を結合が受信した場合にのみ true になります。

ストリーム入力が1つのデータ加工演算子の場合、ウォーターマークは、各並列インスタンスで確認されたイベント時刻の最大値のうちの最小値です。ストリーム入力が2つ以上のデータ加工の場合、ウォーターマークは、入力のウォーターマークの最小値です。

リプレイでは似た出力になりますが、多少異なる場合もあります。これは、リプレイ時に異なるパーティションキーが異なる Flink 並列インスタンスに割り当てられる可能性があることと、入力が複数ある演算子では、上流で入力が異なる速度で処理されていた可能性があることによります。Flink の処理時刻はサポートされておらず、推奨もされません。リプレイ時に結果が大きく異なり、直感に反するものになる可能性があるためです。

並列インスタンスがレコードを受信していない場合、そのインスタンスはウォーターマークを生成しません。これにより、データ加工演算子全体のウォーターマークが進まなくなります。これは多くの場合、キーの分散が不適切であることを示しており、状態の有効期限やウィンドウに重要な影響を及ぼします。これを解決するには、Pipeline Builder でストリームのアイドル状態を設定します。

並列インスタンスが、設定された処理時間にわたってレコードを受信しない場合、そのウォーターマークはデータ加工演算子全体のウォーターマークの計算で考慮されなくなります。この設定によってウォーターマークが停滞する問題を解決できますが、タイムアウトを短くしすぎると、低速なインスタンスが誤ってアイドル状態に指定される可能性があります。その場合、高速なインスタンスがデータ加工演算子全体のウォーターマークを進めるため、破棄されるレコードが増える可能性があります。

状態の有効期限

大きな状態を保存すると、パフォーマンスのボトルネックが生じ、スループットや低遅延に悪影響を及ぼす可能性があるため、Pipeline Builder ではユーザーが状態のサイズを制限する必要があります。

通常、状態はユーザーが指定したキャッシュ時間の有効期限によって制限されます。キャッシュ時間パラメーターが必要なステートフルなトランスフォームでは、通常、ウォーターマークが、そのキーで最後に観測されたイベント時刻に有効期限を加えた時刻を超えるまで、キーごとに状態が cache に保存されます。

ウォーターマークが停滞すると、ステートフルなトランスフォームは状態を速やかに削除しません。これにより、予期しない出力が発生したり、状態が際限なく増大したりする可能性があります。

ウィンドウとトリガー

ウィンドウ集計トランスフォームでは、行とその状態をまとめてグループ化するための方式である*ウィンドウと、集計が出力を生成するタイミングを決める方式であるトリガー*を設定できます。

ウィンドウ

現在サポートされているウィンドウは次のとおりです。

  • タンブリングイベント時間: 時間を、重複しない連続した固定長の区間に分割します。同じキーを持ち、イベント時刻が同じ区間内にある行がまとめてグループ化されます。たとえば、同じ日付で、イベント時刻が同じ1時間の区間内にあるすべての行をまとめてグループ化できます。
  • カウント: ユーザーが指定した件数 n に基づいて、キーごとに、そのキーを持つ最新の n 行をまとめてグループ化します。
  • セッション: 同じセッションに属する行をまとめてグループ化します。キーが同じで、そのキーを持つ行のイベント時刻の間隔がユーザー指定のセッションギャップを超えなければ、行は同じセッションに属します。たとえば、ストリーミングプラットフォームでのユーザー操作に関するデータを含むデータセットでは、ユーザーが操作を中断するまで、1人のユーザーのワークフローに関するすべての行をまとめてグループ化できます。

時間に依存するウィンドウ(タンブリングイベント時間ウィンドウやセッションウィンドウなど)は、ウォーターマークが十分に進むと、最終的に閉じられます。

  • 許容遅延が設定されていないか、ゼロの場合、ウィンドウはウォーターマークがウィンドウの終了時刻を過ぎるまで開いたままになります。その時点でウィンドウが閉じられ、出力が生成される場合があり、状態が削除されます。
  • 許容遅延が指定されている場合、ウィンドウはウォーターマークがウィンドウの終了時刻に許容遅延を加えた時刻を過ぎるまで開いたままになります。これにより、ウォーターマークがウィンドウの終了時刻を過ぎていても、遅れて到着したレコードや順序どおりに到着しなかったレコードをウィンドウに含めることができます。
  • ウォーターマークがウィンドウの終了時刻に許容遅延を加えた時刻を過ぎてから到着した行は、ウィンドウがすでに閉じられ、状態も削除されているため、必ず破棄されます。

時間に依存するウィンドウでは、カスタムトリガーも指定できます。

ウォーターマークが停滞すると、ウィンドウが開いたままになる時間が大幅に長くなります。これにより、行が発行されない、状態が際限なく増大する、トリガーが発火しないといった問題が生じる可能性があります。

トリガー

現在サポートされているトリガーは次のとおりです。

  • ウォーターマーク後トリガー: ウォーターマークがウィンドウの終了時刻を過ぎると、ウィンドウが出力を生成します。ウォーターマークがウィンドウの終了時刻より前の場合と、終了時刻より後の場合(許容遅延によってウィンドウがまだ存続している場合)に使用する、別のカスタムトリガーを指定できます。たとえば、ウィンドウが閉じるまでは出力を生成せず、許容遅延期間内に遅れて到着したレコードについてはすべての出力を確認したい場合に使用できます。
  • 件数トリガー: ユーザーが指定した件数 n に基づいて、キーごとに、受信した行数が n の倍数に達するたびにウィンドウが出力を生成します。
  • ウィンドウ終了トリガー: ウィンドウが閉じられ、状態が削除されるときにのみ出力を生成します。ウィンドウごとに必ず1回だけ、ウィンドウの終了時にのみ発火します。

ステートフルストリーミングのベストプラクティス

大きな状態はパフォーマンスに悪影響を及ぼす可能性があるため、ステートフルなパイプラインを設計する際は、状態の有効期限ポリシーをできるだけ「厳しく」設定することを推奨します。通常、これはキャッシュ時間の有効期限を必要以上に長く設定せず、件数ウィンドウの件数も必要以上に大きく設定しないことを意味します。

大きな状態が必要なパイプラインでは、パフォーマンス(スループット、Checkpoint の所要時間、低遅延を含む)は Flink ジョブの並列度に応じて向上します。並列度はストリーミングパイプラインの設定で編集できます。並列度を大きくすると、データ処理能力と状態の読み取りおよび書き込み速度が向上します。

キー列の値の種類が多すぎたり、行の分布に偏りがあったりすると、ボトルネックやスケーリングの問題が生じる可能性があるため、ステートフルなトランスフォームには適切なキーを選択する必要があります。