7.6 リアルタイムデータとストリーミングデータ
概要と動機
データパイプラインについて知っていることのほとんどは、データが静止していることを前提としています。一日分のレコードを集め、一晩ジョブを実行し、朝に結果を読む。リアルタイムデータとストリーミングデータは、その前提をひっくり返します。完成したデータの山を処理するのではなく、到着するままに終わりのないイベントの流れを処理し、継続的に答えを生み出します。これが、有界で完全なデータセットに対して動作するバッチ処理と、無制限で終わらない流れに対して動作するストリーム処理の違いです。
大きなチームにとって、ストリーミングは、レイテンシがビジネスにとって重要になった瞬間に現れます。一時間遅れて届く不正の判断は無価値です。明日着地するパーソナライズのシグナルは、何もパーソナライズしません。現実に一勤務分遅れる運用のダッシュボードは、それを見ている人々を誤導します。7.2章(データエンジニアリング)は、バッチを既定にし、レイテンシが本当に見返りをもたらす所でのみストリーミングに手を伸ばすべきだと論じており、本章は残りの道のりを進みます。リアルタイムがそのコストに見合うのはいつか、運用予算に火をつけずにそれをどう築くか。ストリーミングは、3.12章(イベント駆動アーキテクチャとメッセージング)のイベント駆動のメッセージングのパターン、3.4章(データアーキテクチャとストレージ)のストレージの選択、9.2章(オブザーバビリティとテレメトリ)のテレメトリの実践に近い所にあります。
企業と政府の設定は、賭け金を上げます。銀行は、カードリーダーがまばたきする間に、不正のために取引をスコアリングします。交通機関は、車両を追跡し、何百万もの乗客のために到着を予測します。給付機関は、すべての決定の監査可能な記録を保ちながら、請求の異常を監視します。これらすべてで、価値はデータが新鮮なうちに行動することから来て、リスクは、間違っている、不完全である、後で再構築できないデータに基づいて行動することから来ます。本章は、その両方について主張を持っています。
主要原則
- レイテンシに明確なビジネス上の価値があるときにだけストリーミングに手を伸ばします。バッチのほうが安く単純です。
- 有界(有限)のデータと無制限(終わりのない)のデータを区別し、それに応じて設計します。
- 到着時刻ではなく、イベント時刻を真実の源として扱い、遅れた、あるいは順序の乱れたデータに備えます。
- ウィンドウとウォーターマークが、無限のストリームから有限の答えを得る方法です。
- 脆い正確に一度の約束より、冪等なシンクを通じた実質的に一度の結果を好みます。
- 状態を持つ処理には、データを失ったり二重に数えたりせずに回復できるよう、チェックポイントが必要です。
- 後付けではなく、初日からバックプレッシャーと再処理を設計します。
- ストリーミングのロジックを観察可能で監査可能に保ちます。静かなストリームは、失敗したバッチより悪い。
推奨事項
築く前にリアルタイムを正当化する
最も重要なストリーミングの決定は、そもそもストリーミングするかどうかです。リアルタイムは、運用の複雑さとコストをおよそ倍にします。実行して止まるジョブを、毎秒健全でなければならないシステムと交換するからです。コミットする前に、新鮮なデータが可能にする決定と、その決定が遅れて届くコストを名指しします。不正のスコアリング、運用のアラート、ライブのパーソナライズは、通常、基準を満たします。人間が一日に二度見るダッシュボードは、計画の会議で「リアルタイム」がどれほど満足に聞こえても、ほとんど満たしません。レイテンシの要件を、秒あるいは分の数字として書き留め、現実と照らしてください。人々がリアルタイムと呼ぶもののかなりは、数分ごとに走るマイクロバッチが、コストのごく一部で十分に仕えます。
処理時刻ではなくイベント時刻を軸に設計する
ストリーミングで最も難しい唯一の考えは、イベントはある瞬間に起こり、別の瞬間に処理されるということです。イベント時刻は、実際にそれが起こった時で、たとえば乗客がカードをタップしたときです。処理時刻は、システムがそれを扱うことに取りかかった時です。これらは絶えずずれます。トンネルで電波を失った電話が三分間のタップを一度にアップロードする、ネットワークの瞬断がメッセージの順序を入れ替える、パーティションが遅れる。処理時刻で計算すると、数字は世界を反映するのではなく、インフラとともに揺れます。この遅れた、順序の乱れた問題が規律の核心で、イベント駆動アーキテクチャのイベントモデリングに直接つながります。すべてのイベントに、源でそのイベント時刻を刻印し、そのタイムスタンプをパイプライン全体を通じて運び、それに対して結果を計算します。
ウィンドウとウォーターマークで有限の答えを得る
無制限のストリームは終わらないので、「イベントを数える」には、それを有界にするまで答えがありません。ウィンドウがその有界化を行います。タンブリングウィンドウは、時間を固定の重ならないバケツ、たとえば毎分に切ります。スライディングウィンドウは重なるので、毎分進む5分のウィンドウが滑らかな移動する数字を与えます。セッションウィンドウは、非活動の間隔で区切られた活動の群れをグループ化し、ユーザーのセッションによく合います。ウィンドウができたら、遅れたデータがまだ届くかもしれないので、ウィンドウがいつ完了したかを決める必要があります。ウォーターマークは、あるイベント時刻までのすべてのイベントをおそらく見たという、システムの推定です。ウォーターマークがウィンドウの終わりを過ぎたら、結果を出します。どれだけ待つかを調整します。ウィンドウを長く開けておけば、レイテンシとメモリを代償に、より多くの遅れを許容でき、早く閉じれば、取り残されたものを落とすリスクがあります。ウィンドウが閉じた後に届いたデータに何が起こるかを、落とすか、記録するか、訂正を出すかを明示的に決めます。
シンクを冪等にし、実質的に一度を好む
配信の保証は単純に聞こえますが、そうではありません。少なくとも一度の配信は、すべてのイベントが処理されることを意味しますが、再試行の後で一部が複数回処理されることがあり、カウントが膨らみえます。正確に一度は理想的に聞こえますが高価で、任意の外部システムにわたって文字どおりには、しばしば不可能です。実用的な目標は実質的に一度です。裏で再試行があっても、観察される結果が、各イベントが一度処理されたかのようであること。そこに至るには、決定的なキーとアップサートを使い、再生されたイベントが重複するのではなく上書きするよう、冪等なシンクを繰り返し書き込んで安全にします。少なくとも一度の配信と冪等な書き込みを組み合わせると、あらゆる所で重い取引の調整に払わずに、正しい結果が得られます。本物の正確に一度の仕組みは、本当にそれを必要とする狭い場所のために取っておきます。
状態を持つ処理が回復できるようにチェックポイントを取る
多くの有用なストリーミングの計算は状態を持ちます。累積のカウント、ストリームをまたぐ結合、重複排除、最近の振る舞いを覚えている不正のモデル。その状態はメモリに住み、プロセスが再起動すると消えます。チェックポイントは、状態とストリームの位置を一緒に定期的にスナップショットするので、クラッシュの後、システムはすべてを再生したりメモリを失ったりせず、一貫した点から再開します。状態のサイズを意図して見積もります。無制限の状態は、本番でストリーミングジョブをメモリ不足にする一般的な方法だからです。もう必要としない状態には、有効期限と生存時間を使い、状態のサイズを第一級の指標として監視します。失敗後の回復時間は本物のサービスレベルの関心事なので、ユーザーより先にテストしてください。
変更データキャプチャで運用データベースからストリーミングする
イベントを発するようには設計されなかったデータベースの変更に、反応したいことがよくあります。変更データキャプチャ(CDC)は、データベースのトランザクションログを読み、すべての挿入、更新、削除を変更イベントのストリームに変えることで、これを解決します。これは、遅く、中間の状態を見逃し、源に負荷をかける、タイマーでテーブルをポーリングするよりはるかに優れています。CDCにより、検索インデックス、キャッシュ、分析ストア、下流のサービスを、記録のシステムと継続的に同期でき、アプリケーションへの侵襲的な変更なしに行えます。変更ストリームを第一級のデータプロダクトとして扱います。そのスキーマにバージョンを付け、意味を文書化し、遅れを観察します。下流のすべてがその遅れを引き継ぐからです。
二つのコードベースを保守するより、ストリーミングファーストのアーキテクチャを好む
古典的なラムダアーキテクチャは、正確で完全な履歴のためのバッチ層を、新鮮で近似的な結果のためのスピード層と並べて動かし、それからそれらをマージします。機能しますが、同じビジネスロジックを二つのシステムで二度書いて保守し、差異を永遠に突き合わせることを強います。カッパアーキテクチャはこれを畳みます。持続的で再生可能なイベントのログを保ち、すべての処理をストリーム処理として動かし、ロジックが変わったときはログを再生して履歴を再処理する。単一のコードベースが保守も推論も劇的に安いので、業界はこのストリーミングファーストの形へ流れてきました。バッチのニーズを、保持されたイベントログに対する再生として表現できるなら、二つのコードベースの税を完全に避けられます。再処理が再構築ではなく巻き戻しの問題になるよう、履歴を保持するログベースのブローカーを使ってください。
ストリームをSQL、マテリアライズドビュー、リアルタイムOLAPとして公開する
ストリーミングを必要とする全員が、低レベルのストリーム処理のコードを書かなければならないわけではありません。ストリーミングSQLにより、アナリストとエンジニアは、すでに知っている言語でウィンドウ、結合、集計を表現でき、その結果をマテリアライズドビューとして継続的に最新に保ちます。新鮮なデータに対する低レイテンシの分析クエリには、リアルタイムのオンライン分析処理(OLAP)ストアがストリームを取り込み、スライスアンドダイスのクエリにミリ秒で答え、それが本当にライブの運用のダッシュボードを動かすものです。機能と実験への素早いフィードバックが目標のときは、7.4章(プロダクト分析と実験)のプロダクト分析の実践と組みます。これらのより高レベルのツールを合う所で選び、手書きのストリームプロセッサは、それらが表現できないロジックのために取っておきます。
最初からバックプレッシャーと再処理を計画する
ストリームは、処理できるより速く届くことがあります。バックプレッシャーは、遅い消費者が、倒れたりデータを静かに落としたりする代わりに、上流に減速するよう合図する仕組みです。パイプラインのすべての段階がそれを尊重することを確認し、消費者の遅れを見出しの指標として監視します。遅れの増大が、競争に負けつつある最も早い警告だからです。再処理は、人々が組み込んでおけばよかったと願うもう一つの能力です。バグを見つけたりルールを変えたりしたとき、履歴を修正されたロジックに通して再生したくなります。それが可能なのは、イベントログが十分な履歴を保持し、シンクが再生を吸収できるほど冪等な場合だけです。両方を初日から設計してください。インシデントの圧力のもとでの後付けは悲惨です。
トレードオフ: 長所と短所
| 選択 | 長所 | 短所 | 最適な場合 |
|---|---|---|---|
| バッチ | 単純、安く、テストとバックフィルが容易 | 高いレイテンシ。実行の間は古い | レポーティング、ほとんどの分析 |
| マイクロバッチ(分) | ほぼリアルタイム。ストリーミングよりはるかに単純 | 本当に即時ではない | 「リアルタイム」のダッシュボード |
| 本物のストリーミング(1秒未満) | 即時の反応。継続的な結果 | 複雑、高価、テストが難しい | 不正、アラート、ライブのパーソナライズ |
| 少なくとも一度 + 冪等なシンク | 正しい結果。手頃。レジリエント | 規律あるキーの設計が必要 | ほとんどのストリーミングのパイプライン |
| 正確に一度の仕組み | 端から端までの強い保証 | 高価。システム間では限られる | 狭い賭け金の高い経路 |
| ラムダ(バッチ + スピード) | 正確な履歴と新鮮なビュー | 保守する二つのコードベース | レガシーの移行 |
| カッパ(ストリーミングファースト) | 一つのコードベース、再生可能 | 保持された持続的なログが必要 | 新しいストリーミングのプラットフォーム |
中心的な緊張は、レイテンシ対複雑さです。リアルタイムへの一歩ごとに、運用の負担、テストの難しさ、お金のコストがかかり、見返りは線形ではありません。毎日から数分ごとへは安く、しばしば十分で、分から1秒未満へは、費用が集中する所です。技術ではなく決定を値付けすることで、緊張を解決してください。鮮度がどんな行動を可能にし、遅れが何を犠牲にするかを問い、その行動が正当化する分だけレイテンシの削減を買います。ストリーミングが必要なときは、少なくとも一度の配信と冪等なシンクとストリーミングファーストのログに頼ってください。その組み合わせは、最も重い保証なしに、正しさと再生可能性を与えるからです。
チームで議論すべき問い
リアルタイムのデータは、私たちにどんな決定を実際に可能にし、そのデータが即時ではなく一分遅れて届くと、何を犠牲にしますか。 これは、ストリーミングが運用のコストと複雑さをバッチに比べておよそ倍にするので、すべてのストリーミングのプロジェクトをゲートすべき問いです。大きなチームは、人間が一日に二度確認するダッシュボードに仕えるリアルタイムのプラットフォームを築くのに四半期を費やしえ、それはお金に火をつけることです。データが駆動する具体的な行動、不正な取引のブロック、オペレーターの呼び出し、ユーザーが見るものの変更のいずれであれ、を持ち込み、それぞれについてレイテンシのコストに数字を付けてください。誠実な答えが、5分のマイクロバッチでニーズに足りる、というものなら、それは隠すのではなく祝う価値のある発見です。答えは、本物のストリーミングを築くか、マイクロバッチで済ませるか、バッチに留まるかを直接変えるべきです。
遅れた、順序の乱れたイベントをどう扱い、ウィンドウが閉じた後に届いたデータには何が起こりますか。 遅れた、順序の乱れたデータはストリーミングの難しい部分で、この問いを飛ばすチームは、数字が突き合わなくなる本番でそれを発見します。相反する圧力はレイテンシと正しさです。取り残されたものを捉えるためにウィンドウを長く開けておけば、あらゆる結果を遅らせてより多くのメモリを消費し、早く閉じれば、本物のデータを静かに落とします。データが実際にどれだけ遅れて届くかの証拠を持ち込んでください。源にわたるイベント時刻と処理時刻のギャップとして測定します。トンネルの中のモバイルの源は、サーバーサイドのイベントとは非常に異なる振る舞いをするからです。遅れたデータを落とすか、記録するか、訂正を引き起こすかを明示的に決め、下流の全員がどれかを知るようにしてください。数字が擁護できなければならない政府の文脈では、遅れたイベントを静かに落とすことはコンプライアンスの問題になりえるので、方針は意図され、文書化される必要があります。
シンクは履歴を安全に再生できるほど冪等ですか。そしてイベントログは再生を可能にするほど保持していますか。 再処理は、チームが最も組み込んでおけばよかったと願い、最も組み込まなかった能力で、二つが一緒に機能することに依存します。重複せずに再生されたイベントを吸収する冪等なシンクと、再生元となる十分な履歴を保持する持続的なログ。両方がなければ、ロジックのバグの修正は、影響を受けた期間をきれいに再計算できないことを意味し、圧力のもとで手作業で数字にパッチを当てることになります。現在の保持の窓と具体的なテストを持ち込んでください。先四半期の本物のバグを一つ選び、修正されたロジックを影響を受けたデータに再生できたかを問ってください。これに逆らう引力はコストです。履歴を保持して冪等な書き込みを設計するには、ストレージと前もっての規律がかかるからです。しかし代替案は、インシデントの最中という最悪の瞬間に表面化するので、答えは、必要になる前にどれだけ再生可能性に投資するかを形づくります。
ストリーミングジョブがクラッシュしたとき、どれだけ速く回復しなければならず、どれだけの状態を保持してよく、本番の負荷のもとでの回復を実際に計時しましたか。 死んだバッチジョブは明日再実行できますが、死んだ常時稼働のストリームは進行中の障害で、累積のカウント、結合、不正のモデルを保持する状態を持つジョブは、数分のメモリを失ったり、再起動後に状態を再読み込みするのに長くかかったりしえます。大きなチームにとって、ここは地味な詳細が実際の可用性を静かに設定する所です。無制限の状態はジョブがメモリ不足になるまで育ち、遅いチェックポイントの復元は、10秒の瞬断を10分のものに変えます。相反する圧力は鮮度対安全です。より頻繁なチェックポイントは回復を短くしますがオーバーヘッドを加え、寛大な状態の保持は精度を改善しますがメモリ枯渇のリスクがあります。具体的な目標復旧時間、現在の状態のサイズとその成長曲線、チェックポイントの間隔、希望的な見積もりではなく本物のフェイルオーバー訓練の結果を持ち込んでください。ストリームが不正のスコアリングや公共の安全のフィードを支える企業と政府の設定では、テストされていない回復の経路は、測定せずに受け入れた運用上のリスクなので、訓練を、あれば良いものではなく要件として扱ってください。
一つのストリーミングファーストのコードベースを動かしていますか。それとも別々のバッチ層とスピード層で、その二つを突き合わせ続けるのに実際いくらかかりますか。 正確な履歴のためのバッチ層と新鮮な結果のためのスピード層というラムダのパターンは、同じビジネスロジックを二つのシステムで二度書き、その答えを永遠に突き合わせることを強い、一方、ストリーミングファースト(カッパ)の形は、持続的で再生可能なログを保ち、すべての処理をストリーム処理として動かします。大きな組織にとって、重複したロジックは、ずれと争いのある数字が育つ所です。ルールが一つの層で変わり他では変わらず、エンジニアが二つがなぜ食い違うかを説明するのに本物の時間を費やすからです。両方を保つことへの引力は、慣性と実証済みのバッチ層の心地よさなので、それを保守の税と誠実に量ってください。現在両方の場所で実行している計算の一覧、二つの層が食い違うことに起因するインシデント、イベントログがバッチのニーズを再生として表現できるほど十分な履歴を保持しているかの評価を持ち込んでください。政府と監査される企業の文脈では、同じ期間について異なる数字を報告しえる二つの層は、それ自体がコンプライアンスの負債です。どの数字が権威あるもので、なぜかを言えなければならないからです。
午前3時にこの常時稼働のシステムが壊れたとき、誰が運用し、オンコールの負荷と必要な専門技能に予算を付けましたか。それともバッチ型の人員配置を想定していますか。 ストリーミングは、コストを構築から運用へ移します。システムは毎秒健全でなければならず、それは本物のオンコールの対応、イベント時刻、ウォーターマーク、状態、配信のセマンティクスに通じたエンジニア、実行して止まるジョブよりも難しいテストを意味します。チームは、その能力に基づいてストリーミングのプラットフォームを日常的に承認し、それを生かし続ける人々に決して資金を出さないので、プラットフォームは劣化し、信頼が侵食されます。トレードオフはスコープ対持続可能性です。追加のリアルタイムのパイプラインは、誰かを呼び出しうるもう一つのものなので、問いは、それが買うレイテンシが恒久的な運用のコミットメントを正当化するかです。本番の各ストリームを誰が所有するか、現在のオンコールのローテーションとその余裕、イベント時刻の専門知識が実際にどこにあるか(採用、パートナー、マネージドサービスのどれか)についての誠実な目録を持ち込んでください。公的機関や大企業では、調達と採用のリードタイムとマネージドサービスの選択肢を加えてください。採用も維持もできない稀少な人材に依存するリアルタイムのプラットフォームは、障害の起きやすいシステムを人員不足で運用する計画だからです。
セクター別の視点
スタートアップ。 ストリーミングはめったに最初の動きではなく、重いプラットフォームを立ち上げることは、小さなチームを沈めえます。中核の価値に触れるシグナルを一つ選び、イベントを保持されたログベースのブローカーの上に置き、少なくとも一度の再試行が決して二重に数えないよう、キー付きで冪等なシンクを伴う軽量なプロセッサを動かしてください。修正されたロジックに通して再生できるよう数日の履歴を保ち、自前のクラスターを運用するよりマネージドなストリーミングサービスを好みます。最も乏しい資源はエンジニアリングの注意だからです。
小規模事業者。 おそらくストリーミングの専門家はおらず、常時稼働のインフラストラクチャを動かす意欲もないので、リアルタイムを、人員を置くシステムではなく、すでに使っているツールの中で買うものとして扱ってください。ニーズを数字を伴うレイテンシの問いとして枠づけると、ほとんどの場合、数分ごとに更新されるマイクロバッチが、コストとリスクのごく一部でそれを満たします。リアルタイムの機能が遅れについて透明で、フォールバックしやすいベンダーを選び、カスタムのストリーミングは、新鮮なデータが収益や安全を直接駆動するまれなケースのために取っておいてください。
大企業。 問題は、多くのチームにわたる一貫性とコストです。共有のログベースのプラットフォーム、標準のイベント時刻と遅れたデータの方針、グループが脆いパイプラインを再発明するのをやめるための冪等なシンク。常時稼働の運用とオンコールの負担を明示的に予算化し、重複したバッチのコードベースを避けるためにストリーミングファーストのログに標準化し、ストリームを、散らばったオーダーメイドのジョブではなく、所有者、スキーマのバージョニング、監視される遅れを備えた統治されたデータプロダクトとして管理します。レイテンシ、回復時間、ストリームごとのコストをポートフォリオの指標として追跡してください。
政府。 監査可能性と公的な説明責任があらゆる選択を形づくります。監督機関に報告される数字、乗客数、給付の異常、不正の決定が正確に再構築できるよう、処理されたすべてのイベントを持続的なログに保持し、遅れたデータの方針を、イベントを静かに落とすのではなく、明示的に文書化してください。調達は、データの可搬性と、マネージドサービスの配信と保持の保証の開示を求めるべきで、ルールの変更後の再表明は、誰もたどれない手作業のパッチではなく、修正されたロジックを通した擁護できる再生であるべきです。
事例
スタートアップ。 消費者向けアプリが、ユーザーにライブのアクティビティフィードを見せ、疑わしいログインが起きるたびに旗を立てたいと考えています。チームは重いストリーミングのプラットフォームを立ち上げることに抵抗します。イベントを保持されたログベースのブローカーに置き、ログインのリスクのロジックに軽量なストリームプロセッサを動かし、アクティビティフィードを動かすリアルタイムのOLAPストアに供給します。すべてのシンクはキー付きで冪等なので、少なくとも一度の再試行が二重に数えることは決してありません。後でリスクのルールのバグを見つけたとき、一週間の履歴を保っていて、二つ目のバッチのコードベースを必要としなかったので、修正されたロジックに通してログを一晩で再生するだけです。
大企業。 小売銀行は、直近のアカウントの振る舞いの状態を持つモデルに対してライブの取引ストリームを結合しながら、承認の窓の中ですべてのカード取引を不正のためにスコアリングします。チェックポイントにより、スコアリングのサービスは、ノードの障害から、直近数分の記憶を失わずに数秒で回復できます。別に、変更データキャプチャが、勘定系のデータベースからの更新を検索インデックスとパーソナライズのサービスにストリーミングし、ポーリングなしに両方を新鮮に保ちます。運用のダッシュボードはリアルタイムのOLAPストアから読むので、リスクと運用のチームは事業が動くままに見て、パイプライン全体が9.2章で述べられた遅れとスループットのテレメトリを発します。
政府。 大都市の交通局は、到着を予測し混雑を監視するために、車両の位置と運賃のタップをリアルタイムで取り込み、公共のアプリと運用センターの両方に供給します。トンネルの乗客が遅れたバーストでタップをアップロードするので、チームはイベント時刻で、観察された遅れに調整したウォーターマークとともに乗客数を計算し、ウィンドウが閉じた後に届いたイベントを静かに落とすのではなく記録します。処理されたすべてのイベントは監査可能なログに保持されるので、監督機関に報告される乗客数は正確に再構築できます。運賃のルールが変わったとき、影響を受けた期間を修正されたロジックに通して再生し、擁護できる再表明を作ります。
ビジネスケース: 動機、ROI、TCO
リアルタイムのデータの見返りは、行動がまだ重要なうちに行動することから来ます。承認の間に捉えられた不正は、夜間のバッチが報告するだけの損失を防ぎます。セッション内で応答するパーソナライズは、明日のレコメンデーションにはできない形でコンバージョンを上げます。現在を反映する運用の監視は、小さな問題が障害や公的なインシデントになる前に介入できるようにします。いずれの場合も、価値は今行動することと後で行動することの差で、論拠を示すとき定量化すべきはその差です。
総所有コストはバッチより高く、それについて誠実であることがあなたの信頼性を守ります。常時稼働のインフラストラクチャ、イベント時刻、ウォーターマーク、状態、配信のセマンティクスを理解するエンジニア、実行して止まるのではなく継続的に健全でなければならないシステムのより難しいテストとオンコールの負担に払います。保持されたログ上のストリーミングファーストのアーキテクチャは、重複したバッチのコードベースを省くことで継続的なコストを下げ、冪等なシンクを伴う少なくとも一度の選択は、端から端までの正確に一度の仕組みの費用を避けます。最も高価な間違いは、マイクロバッチやバッチで足りる所でリアルタイムを築くことなので、最も強いコストの論拠は、しばしばストリーミングしないという決定です。リーダーシップへの訴えは、特定のレイテンシに敏感な決定と、その測定可能な見返りを軸に枠づけ、バッチに留まることが価値を失わずにお金を節約する所についても同じくらい明確にしてください。
アンチパターンと落とし穴
- 数分ごとのマイクロバッチでニーズを満たせるのに、威信のためにストリーミングを築くこと。
- 処理時刻で計算し、数字が世界ではなくインフラとともに揺れること。
- 本番で突き合わせが失敗するまで、遅れた、順序の乱れたデータを無視すること。
- 冪等なシンクを伴う少なくとも一度ではなく、あらゆる所で文字どおりの正確に一度を追うこと。
- 有効期限のない無制限の状態が、ジョブがメモリ不足になるまで静かに育つこと。
- チェックポイントがなく、再起動で状態を失うか、完全な再生を強いられること。
- 変更データキャプチャを使う代わりに、タイマーで運用データベースをポーリングすること。
- 重複してずれていくロジックを持つ、ラムダのバッチ層とスピード層を保守すること。
- バグを見つけたときに履歴を再生するには短すぎる保持の窓。
- 遅れ、スループット、鮮度の指標がなく、静かに失敗するストリーム。
成熟度モデル
- レベル1、開始: すべてがバッチか、少数の手作りのストリーミングジョブが監視なしに反応的に走ります。数字は処理時刻で計算され、遅れたデータは無視され、再起動で状態を失います。バグを直すために履歴を再生できる人はおらず、問題は下流の数字が突き合わなくなったときに発見されます。
- レベル2、発展: ログベースのブローカー上でチェックポイントを伴う中核のストリーミングのパイプラインを動かし、イベント時刻と処理時刻を区別して基本的なウィンドウを使うチームもあります。実践はチームごとに一貫しません。配信は少なくとも一度だがすべてのシンクが冪等ではなく、遅れたデータの扱いは即興で、遅れは非公式に観察され、アラートは出ません。
- レベル3、標準化: イベント時刻、ウォーターマーク、明示的な遅れたデータの方針が文書化され、組織全体で適用されます。シンクは実質的に一度の結果のために冪等で、状態には有効期限があり、変更データキャプチャが慣習として下流のシステムに供給します。保持されたログが再生をサポートし、遅れ、スループット、鮮度は、チームごとの習慣ではなく、組織全体の標準としてアラート付きで監視されます。
- レベル4、管理: ストリーミングの資産群が、ベースラインに対して測定され、制御されます。各パイプラインは、端から端までのレイテンシ、消費者の遅れ、回復時間、イベント時刻のずれ、遅れたイベントの率、状態のサイズ、百万イベントあたりのコストのサービスレベル目標を持ち、すべて合意された目標に対して追跡され、退行にアラートを出します。回復は想定されるのではなく訓練され計時され、バックプレッシャーの余裕と状態の成長は容量のシグナルとして観察され、新しいストリームは本番に進む前にこれらの指標を満たさなければなりません。
- レベル5、オーケストレーション: ストリーミングファーストのアーキテクチャが、一つの再生可能なログから新鮮なニーズと履歴のニーズの両方に仕え、ストリーミングSQL、マテリアライズドビュー、リアルタイムのOLAPが、新鮮なデータを広くアクセス可能にします。再処理は日常的でテストされ、プラットフォームは測定された負荷とコストに対して自動スケールして再均衡し、ストリームは証拠に基づいて退役、再スコープ、置き換えられます。ストリーミングはビジネスとリスクの計画と統合され、負荷とコストの状況が移るにつれて、すべてのストリームが端から端まで観察可能で監査可能です。
議論のためのアイデア
- スタックのどこで「リアルタイム」が実際にそのコストに見合い、どこで吟味されない願いですか。
- 源にわたるイベント時刻と処理時刻のギャップはどれくらい大きく、それを測っていますか。
- ラムダのバッチとスピードの構成を単一のストリーミングファーストのコードベースに畳めますか。そして何がそれを妨げますか。
- どのシンクが本当に冪等で、今日、修正されたロジックに通して先四半期のデータを安全に再生できますか。
- ウィンドウが閉じた後に届くデータの方針は何で、下流の全員がそれを知っていますか。
- 変更データキャプチャは、検索、キャッシュ、分析を同期させる方法をどう変えますか。
要点
- レイテンシに敏感な決定がそれを正当化するときにだけストリーミングに手を伸ばします。バッチとマイクロバッチがより安い既定です。
- イベント時刻で計算し、遅れた、順序の乱れたデータを、ウィンドウとウォーターマークで扱う中核の問題として扱います。
- あらゆる所での文字どおりの正確に一度より、実質的に一度の結果のために、冪等なシンクを伴う少なくとも一度の配信を好みます。
- 状態を持つ処理にチェックポイントを取り、状態を有界にし、消費者の遅れを見出しの指標として監視します。
- ポーリングの代わりに、変更データキャプチャを使って運用データベースからストリーミングします。
- 二つのコードベースを保守するより、保持された再生可能なログ上のストリーミングファーストのアーキテクチャを好みます。
- ストリーミングSQL、マテリアライズドビュー、リアルタイムのOLAPでストリームを公開し、すべてのストリームを観察可能で監査可能に保ちます。
参考文献とさらなる読み物
- Tyler Akidau, Slava Chernyak, and Reuven Lax, “Streaming Systems.”
- Martin Kleppmann, “Designing Data-Intensive Applications.”
- Nathan Marz and James Warren, “Big Data” (Lambda architecture).
- Jay Kreps, “Questioning the Lambda Architecture” (O’Reilly Radar).
- Fabian Hueske and Vasiliki Kalavri, “Stream Processing with Apache Flink.”
- Ben Stopford, “Designing Event-Driven Systems.”
- Tyler Akidau and colleagues, “The Dataflow Model” (VLDB paper on windowing and watermarks).