Andrew Tanによる
以前は十分だったスケジュールされた実行
過去10年間のほとんどで、分析Workflowは次のように見えました:
データをソースから抽出します。それをウェアハウスにロードします。dbtでSQLモデルを書きます。それらを数時間ごとに実行するようにスケジュールします。その上にダッシュボードを構築します。ビジネスがなぜ数値がずれているのかを尋ねると、フレッシュネスバッジを確認し、ジョブが6時間前に失敗したことを発見します。
完璧ではありませんでしたが、予測可能でした。分析とエンジニアリングの間の契約は明確でした:モデルはウェアハウスで、スケジュールに基づいて、SQLで構築されます。dbtはその契約をエレガントにしました。バージョン管理されたモデル、自動化されたテスト、依存関係グラフ、コードから生成されたドキュメント。バッチ分析にとって、それは真の飛躍でした。
その後、ビジネスはより新鮮なデータを求め始めました。翌日ではなく、新鮮な。毎時ではなく、新鮮な。イベント駆動の新鮮さ。顧客がプロフィールを更新し、推奨エンジンが彼らが離れる前にそれを知るような新鮮さ。
そして突然、スケジュールされた実行では不十分になりました。
なぜdbtの成功がランタイムのギャップを生んだのか
dbtの核心的な洞察は、モデルロジックを実行インフラストラクチャから分離することでした。あなたはSQLを書きます。dbtはDAG、テスト、ドキュメント、マテリアライゼーション戦略を処理します。実際の計算は、あなたが指示した場所で実行されます — Snowflake、BigQuery、Redshift、Databricks。dbtは気にしません。それはSQLモデルのためのコンパイラでありオーケストレーターであり、ランタイムではありません。
その分離はバッチワークロードにとっては天才的でした。それは分析エンジニアがクラスター管理、パーティショニング戦略、メモリチューニングを管理することなくモデリングロジックを所有できることを意味しました。ウェアハウスがそれらすべてを処理しました。分析エンジニアはセマンティクスに集中しました:このモデルは何を意味するのか、どのようにテストされているのか、誰がそれに依存しているのか。
しかし、その分離は重要なことを仮定しています:実行レイヤーがあなたが投げかけるワークロードを処理できること。そして、現代のウェアハウスに対するスケジュールされたバッチジョブにとって、それは真実です。ストリーミングデータに対するリアルタイム処理にとって、それはそうではありません。
ギャップはSQLにはありません。dbtのSQLは依然としてSQLです。ギャップは、イベントとモデル結果の間で起こるすべてのことにあります:
- フレッシュネスSLA。dbtモデルが15分ごとに実行されても、依然として15分遅れています。ストリーミングコンテキストでは、15分は小さなコスチュームを着たバッチジョブです。
- 順序外のイベント。ストリーミングデータは遅れて到着したり、重複したり、順序が乱れたりします。順序付けされた入力を仮定するSQLモデルは、警告なしに誤った結果を生成します。
- 状態を持つ操作。ウィンドウ関数、セッション化、重複排除 — これらはイベント間で状態を維持する必要があります。dbtのマテリアライゼーション戦略(テーブル、インクリメンタル、エフェメラル)は、イベント時間の状態管理のために設計されていませんでした。
- オペレーショナル結合。注文のストリームをリアルタイムでゆっくり変化するディメンションテーブルに結合することは、
customer_idで2つのウェアハウステーブルを結合することとは異なる問題です。ディメンションはクエリの途中で変わるかもしれません。ストリームは待ちません。
これらはエッジケースではありません。それらはストリーム処理の定義的な特徴です。そして、dbtは設計上、それらすべてを実行レイヤーに委任しています。
ベンダーが実際に販売しているもの
昨年データカンファレンスに参加したことがあるなら、マーケティングのシフトを見たことがあるでしょう。ConfluentはFlink上でのdbtネイティブモデリングを推進しています。FivetranはAI支援のモデル生成を備えたdbt Wizardを発表しました。SnowflakeはDynamic Tablesを発表し、BigQueryはMaterialized Viewsを持ち、DatabricksはDelta Live Tablesに注力しています。
メッセージは一貫しています:あなたはdbt Workflowを維持し、ただ...リアルタイムにすることができます。SQLは同じままです。DAGも同じままです。変わるのは速度だけです。
これはおおよそ半分は真実です。
はい、ストリーミングデータ上でSQLを実行できます。Flink SQL、Spark Structured Streaming、そして増え続けるストリームプロセッサは、見慣れたSQL構文をサポートしています。はい、それらのSQLファイルをバージョン管理し、DAGを構築できます。いくつかのツールはdbtスタイルのテストとドキュメントもサポートしています。
彼らが伝えないのは、問題のもう半分 — ランタイムの半分 — は消えないということです。それは形を変えるだけです。
スケジュールされたバッチから連続ストリーミングに移行すると、dbtが解決する必要のなかった新しい一連の懸念を引き継ぎます:
| バッチの仮定 | ストリーミングの現実 |
|---|---|
| ジョブが開始する時点でデータは完全である | データは決して完全ではなく、遅延到着が通常である |
| 失敗は実行の最後に捕捉される | 失敗はイベントごとに検出され、処理されなければならない |
| スキーマの変更は実行間で発生する | スキーマの変更はストリームの途中で発生する |
| 再処理はジョブを再実行することを意味する | 再処理はログを巻き戻し、イベントを再生することを意味する |
| コストはデータ量に比例する | コストはインフラストラクチャの稼働時間に比例する |
SQLは同じように見えるかもしれません。しかし、それを実行するシステムは、まったく異なる問題のセットを解決しています。
誰も答えたくない所有権の問題
ここで組織的な問題が発生します。dbtは明確な境界を作りました:分析エンジニアはモデルを所有し、プラットフォームエンジニアはウェアハウスを所有します。モデルが失敗した場合、それは通常SQLの問題かデータ品質の問題です。ウェアハウスが遅い場合、それはプラットフォームの問題です。関心の分離は、チームの分離にきれいに対応しました。
ストリーミングはその境界を壊します。
リアルタイムパイプラインが失敗した場合、それはSQLの問題かインフラストラクチャの問題でしょうか?イベントが順序外に到着する場合、分析エンジニアがウィンドウロジックを書き直すのか、プラットフォームエンジニアがストリームプロセッサのウォーターマーク設定を再構成するのか?遅延到着イベントがマテリアライズドビューを破損した場合、それを修正するのはSQLを書いた人か、チェックポイントを管理する人か?
私はこれを3つの方法で処理するチームを見てきましたが、そのうち1つだけが機能します:
オプション1:分析エンジニアがストリーム処理を学ぶ。彼らはチェックポイント、backpressure、イベント時間と処理時間の違い、そして正確に一度のセマンティクスに精通します。これは小規模なチームでシニアな人々がいる場合に機能しますが、スケールしません。
オプション2:プラットフォームエンジニアがKafka以降のすべてを所有する。分析エンジニアはSQL仕様を書き、プラットフォームチームがそれをFlinkやSparkで実装します。これにより分離は維持されますが、翻訳レイヤーが作成されます。SQL仕様はイベント時間の動作について曖昧です。プラットフォームチームは仮定をします。その仮定は6か月後にバグになります。
オプション3:モデリングロジックをランタイムの責任から分離する。分析エンジニアはモデルが何を意味するか — ビジネスロジック、テスト、セマンティクスを所有します。プラットフォームエンジニアはそれがどのように実行されるか — 実行エンジン、状態バックエンド、障害回復を所有します。両方の側が契約に同意します:モデルは順序付けされた入力を期待し、ランタイムはその契約を保証するか、明確なエラーを表示します。
この3番目のオプションは設定が難しいです。両方のチームがインターフェースと障害モードに事前に同意する必要があります。しかし、これは分析エンジニアを分散システムの専門家に変えたり、プラットフォームチームを心を読む人にしたりすることなくスケールする唯一の方法です。
"リアルタイムdbt"が実際に必要なもの
あなたのチームがdbtスタイルのモデルをリアルタイムパイプラインに移行することを真剣に考えている場合、実際に必要なものは — マーケティングデッキにあるものではありません:
イベント時間を理解するランタイム。処理時間だけでなく。イベント時間。「いつこれが到着したのか」と「いつこれが起こったのか」の違いは、正しい結果とダッシュボード上で問題ないように見える微妙に誤った結果の違いです。
明示的な状態管理。ウィンドウ化、重複排除、セッション化はすべて状態を必要とします。その状態はチェックポイントされ、復旧可能で、デバッグのためにクエリ可能である必要があります。システムが先週の火曜日の午後2時47分に何を考えていたかを検査できない場合、プロダクションのインシデントをデバッグすることはできません。
厳しいスキーマ進化。列を追加するのは簡単です。型の変更、名前の変更、列の意味のセマンティックシフトを処理するのは難しいです。あなたのパイプラインはこれらの変更を検出し、それらが安全かどうかを判断し、適応するか、明確なエラーで停止する必要があります。
リプレイとバックフィルを一級の操作として扱う。バッチでは、再処理はジョブを再実行することです。ストリーミングでは、ログを巻き戻し、同じロジックを通してイベントを再生することです。あなたのリアルタイムパイプラインが正確にリプレイできない場合、バグから回復することはできず、変更を履歴データに対してテストすることはできず、コンプライアンスを証明することもできません。
ワークロードによるコストの可視性。バッチコストは理解しやすいです:このジョブはこの量のデータを処理し、この時間を要しました。ストリーミングコストは継続的です:ジョブは常に実行され、常にリソースを消費しており、コストはビジネスの成果と明確に相関しません。インフラストラクチャの支出をパイプラインの動作に結びつけるテレメトリが必要です。
これらはSQLの機能ではありません。それらはランタイムの機能です。そして、それらは10分間実行されるデモと10か月間実行されるシステムの違いです。

正直な売り込み
dbtは分析チームの働き方を変えました。それはソフトウェアエンジニアリングの実践 — バージョン管理、テスト、ドキュメント — を、非常に必要としていた分野にもたらしました。その貢献は現実であり、持続的です。
しかし、dbtはデータがバッチで移動し、ウェアハウスが重力の中心であり、「新鮮」とは「この時間に更新された」ことを意味する世界のために構築されました。業界は、データが継続的に移動し、モデルがフライト中に適用され、「新鮮」とは「このミリ秒に更新された」ことを意味する世界に向かっています。
それはdbtを時代遅れにするものではありません。それは、成長するユースケースのセットに対してdbtが不完全であることを意味します。
「リアルタイムdbt」を販売しているベンダーは、実際の需要に応えています。しかし、彼らが販売しているのは通常、ストリームプロセッサ上のSQLインターフェースであり、ストリーム処理がもたらすランタイムの問題への解決策ではありません。SQLは簡単な部分です。状態管理、障害回復、スキーマ進化、運用の可観測性が難しい部分です。そして、それらの難しい部分は、SQLが見慣れたものであるからといって消えるわけではありません。
これらのツールを評価する際には、「私のdbtモデルをより速く実行できますか?」と尋ねるのではなく、「ノードがウィンドウの途中で再起動したときに何が起こりますか?」、「昨日変更したモデルを通して先週の火曜日のデータを再生するにはどうすればよいですか?」、「午前2時にスキーマが変更されたときにシステムは何をしますか?」と尋ねてください。
これらの質問への回答は、あなたがより速いバッチツールを購入しているのか、実際のストリーミングランタイムを購入しているのかを教えてくれます。
私たちの役割
layline.ioでは、バッチとストリーミングの両方を同じパイプラインで処理するランタイムを構築しました。SQLの表面をそれぞれにかぶせた2つの別々のシステムではありません。スケジュールされたウェアハウスロードとリアルタイムイベント処理を、ツール、コンテキスト、またはメンタルモデルを切り替えることなく同じチームが構築できる1つのシステムです。
分析エンジニアはモデルのセマンティクスの所有権を維持します。プラットフォームチームは実行インフラストラクチャの所有権を維持します。しかし、両方の側が同じ環境で、同じ可観測性を持ち、リプレイ、状態回復、スキーマ進化に関する同じ保証を持って作業します。
dbtはモデリングロジックが厳密さに値することを業界に教えました。私たちはそのアイデアを基に構築し、リアルタイムデータが要求するランタイムの厳密さを追加しています。
Andrew Tanは、layline.ioの創設者であり、バッチとリアルタイムのワークロードをスケールで処理するエンタープライズデータ処理インフラストラクチャを構築するシリアルアントレプレナーです。



